Elasticsearch 学习笔记
Elasticsearch 写入时先将每篇文章分词,再反向建立 “单个词汇→包含该词汇的所有文章” 的倒排索引,同时对词汇排序以支撑后续高效查询;搜索时则先借助内存中的 Term Index 前缀定位 + 二分查找快速找到目标词汇对应的文章 ID 列表,再根据 AND/OR 等搜索条件对多个词汇的文章 ID 列表取交集或并集,最后依据筛选后的文章 ID 提取完整内容并返回。
https://www.bilibili.com/video/BV1yb421J7oX
一、顶层设计:为什么需要 Elasticsearch?
1.1 ES 的定位
1.2 传统数据库面临的问题
- 模糊查询性能瓶颈:MySQL 等关系型数据库使用
LIKE %keyword%进行左模糊或全模糊查询时,无法利用索引(Index),会导致全表扫描,在百万级数据量下性能急剧下降。 - 复杂维度筛选困难:当面临海量数据的多维度筛选、统计、标签聚合时,SQL 语句变得极度复杂且执行效率低下。
- 文本相关性缺失:传统数据库只能做“精准匹配”或“简单的包含匹配”,无法计算相关性得分(Score),无法按“搜索结果匹配度”对结果进行排序。
- 分词能力弱:无法处理同义词、纠错、中文分词(如将“由于”和“游泳”区分开)等复杂的自然语言处理需求。
1.3 Elasticsearch 解决的核心问题
- 性能全文检索:基于倒排索引,实现亿级数据毫秒级响应。
- 高可用与横向扩展:天生的分布式架构,通过增加节点即可线性提升存储容量和计算能力。
- 复杂聚合分析:提供强大的聚合(Aggregations)框架,能够替代部分 OLAP 场景,实时统计数据分布。
- 相关性排序:基于 TF-IDF / BM25 算法,提供符合人类直觉的搜索结果排名。
二、核心原理:如何解决问题?
2.1 倒排索引(Inverted Index)
这是 ES 快如闪电的根本原因。
- 正排索引(Forward Index):文档 ID -> 文档内容(类似 MySQL 主键查询)。
- 倒排索引:单词(Term)-> 包含该单词的文档 ID 列表。
- ES 在写入数据时,会通过**分词器(Analyzer)**将文本拆解为单词,建立索引。
- 查询时,直接根据单词找到文档 ID 列表,并通过位运算快速合并结果,无需扫描全表。
2.2 Term Index(词项索引)
- 面临的问题:Term Dictionary(词项字典)数据量极大(千万/亿级),无法全部放入内存,必须存在磁盘。每次查询都去磁盘翻字典,I/O 消耗巨大。
- 解决方案:Term Index 是一个基于词项前缀构建的精简目录树(类似 Trie 树或 FST - Finite State Transducers)。
- 内存驻留:它体积非常小,可以完全加载到 内存(RAM) 中。
- 加速原理:查询时,先在内存中查 Term Index,找到该 Term 在磁盘 Term Dictionary 中的大概位置(Offset 偏移量),然后再去磁盘读取具体内容。
- 类比:Term Dictionary 是字典本身(在磁盘),Term Index 是字典的“拼音首字母索引页”(在内存)。
2.3 存储结构
当倒排索引帮我们找到文档 ID 后,我们还需要获取内容或进行排序。
- Stored Fields(行式存储)
- 用途:用于存储原始文档内容(
_source字段)。 - 特点:行式存储,适合展示完整信息,但不适合聚合分析(读取很多冗余数据)。
- 用途:用于存储原始文档内容(
- Doc Values(列式存储)
- 用途:专用于排序(Sorting)和聚合(Aggregations)。
- 原理:空间换时间。将散落在不同文档中的同一个字段值,集中在一起列式存储。
- 优势:排序或计算平均值时,CPU 可以连续读取内存地址,极大提升效率,且对操作系统文件缓存(OS Cache)非常友好。
2.4 Segment 与 Lucene
- Lucene:ES 的底层核心库。一个 ES 的分片(Shard)本质上就是一个完整的 Lucene 索引。
- Segment(段):
- 最小单元:Lucene 内部由多个 Segment 组成。一个 Segment 包含了倒排索引、Term Index、Stored Fields 等所有结构,是具备完整搜索功能的最小单元。
- 不可变性(Immutable):Segment 一旦生成,不可修改。
- 写入与合并:新增数据会生成新的 Segment;删除数据只是打上
.del标记(逻辑删除)。ES 会在后台自动进行 Segment Merging(段合并),将多个小 Segment 合并为大 Segment,同时物理剔除被标记删除的数据。
三、分布式架构设计
- Cluster(集群):由多个节点组成。
- Node(节点):单个服务实例。
- Index(索引):逻辑上的数据集合(类比 Database)。
- Shard(分片):数据被切分成多个分片存储在不同节点,实现并行计算和存储扩展。
- Replica(副本):分片的备份,用于提高可用性(HA)和读取吞吐量。
3.1 高性能优化:分片(Shard)
- 问题:单个 Index 数据量过大(如 1TB),单机硬盘存不下,且搜索时单线程扫描太慢。
- 解决方案:分片机制。
- 将一个 Index 逻辑拆分为多个 Shard(分片)。
- 每个 Shard 是一个独立的 Lucene 实例。
- 优势:读写压力被分散到多个 Shard 上并行处理,极大提升吞吐量。
3.2 高扩展性优化:多节点(Node)
- 原理:横向扩展(Scale Out)。
- 机制:当数据量增长,分片增多,单机 CPU/内存吃紧时,可以向集群中加入新的机器(Node)。
- 自动平衡:ES 会自动感知新节点,并将部分 Shard 迁移过去,实现负载均衡。
3.3 高可用优化:副本(Replica)
- 角色区分:
- Primary Shard(主分片):负责处理写入请求。
- Replica Shard(副本分片):主分片的完整备份。
- 机制:
- 读写分离:副本分片可以分担搜索(读)请求,提升查询并发量。
- 故障转移(Failover):如果持有主分片的 Node 挂掉,集群会迅速选举一个副本分片升级为新的主分片,保证服务不中断。
3.4 节点角色分化
在大型集群中,让每个节点“各司其职”效率更高。
- Master Node(主节点):集群的大脑。负责索引创建/删除、维护集群状态(Cluster State)、管理节点加入/退出。
- Data Node(数据节点):苦力。负责存储数据(Shard),执行耗费资源的 CRUD 和聚合操作。
- Coordinate Node(协调节点):前台/路由。接收客户端请求,分发给 Data Node,并汇总最终结果(Scatter-Gather)。
3.5 去中心化协调机制
- Raft 算法衍生:ES 内部实现了一套基于 Raft 改进的共识算法(Zen Discovery)。
- 作用:
- 保证集群中所有节点对“集群状态”的认知是一致的。
- 实现 Master 节点的选举。
- 故障检测:节点之间互相 Ping,感知是否有节点掉线。
四、核心流程详解
4.1 写入流程
ES 的写入流程设计权衡了数据安全性与写入高吞吐。
1. 路由与转发
- 客户端向任意节点发送写入请求(该节点暂时成为协调节点)。
- 协调节点使用路由算法确定数据所属的主分片(Primary Shard)位置:
shard = hash(routing) % number_of_primary_shards- 注:
routing默认是文档_id。
- 协调节点将请求转发给持有该主分片的 Data Node。
2. 主分片写入 (Primary Operation)
- 写入内存缓冲区 (Memory Buffer):数据先写入内存 buffer,此时数据不可被搜索。
- 写入 Translog (Transaction Log):同时追加写入 Translog 文件(顺序写磁盘),防止断电丢失数据。
- Refresh (准实时关键步骤):默认每 1 秒,ES 将 buffer 中的数据生成一个新的 Segment 文件(此时建立倒排索引),并清空 buffer。一旦生成 Segment,数据即可被搜索。这就是 ES 被称为“准实时(Near Real-Time, NRT)”的原因。
3. 同步副本 (Replication)
- 主分片写入成功后,并行将请求发送给所有的 副本分片 (Replica Shards)。
- 副本分片执行相同的写入逻辑。
4. 响应客户端
- 当所有在 ISR (In-Sync Replicas) 列表中的副本都反馈写入成功后,主分片向协调节点报告成功。
- 协调节点向客户端返回“写入完成”。
技术深挖:Flush 操作 Translog 不会无限增长。当 Translog 达到阈值或每隔 30 分钟,ES 会触发 Flush 操作:
- 强制执行 Refresh。
- 将所有内存中的 Segment 强制
fsync刷入物理磁盘。- 清空 Translog。 这保证了数据的持久化存储。
4.2 搜索流程(Query Then Fetch)
搜索比写入复杂,因为数据分散在多个分片上,必须通过“两阶段”策略来整合结果,以避免网络带宽的巨大浪费。
阶段一:查询阶段 (Query Phase)
- 请求分发:客户端向协调节点发送搜索请求。协调节点根据请求(是否有 routing 参数)将请求广播到所有相关分片(主分片或副本分片均可,负载均衡)。
- 本地检索:每个分片在本地 Lucene 中执行搜索:
- 利用倒排索引筛选匹配文档。
- 利用 Doc Values 进行排序和打分。
- 关键点:分片仅返回文档 ID、相关性算分 (_score) 和排序值给协调节点,不返回文档的完整内容 (
_source)。
- 全局排序:协调节点收到所有分片返回的轻量级列表(例如每分片前 10 条),在内存中进行全局归并排序,选出最终的 Top N 文档 ID。
阶段二:获取阶段 (Fetch Phase)
- 精确定位:协调节点知道了最终需要哪几个文档,以及它们位于哪个分片。
- 抓取内容:协调节点向相关分片发送
Multi-Get请求,只索取这 Top N 文档的完整内容 (_source/Stored Fields)。 - 返回结果:分片返回文档详情,协调节点拼装最终 JSON,响应给客户端。
性能隐患:深度分页 (Deep Pagination) 如果查询
from=10000, size=10:
- 每个分片都必须查询出前 10010 条记录。
- 假设有 5 个分片,协调节点需要接收
5 * 10010 = 50050条记录的 ID,并在内存中排序,最后只取 10 条。- 后果:内存爆炸,CPU 飙升。
- 对策:避免深分页,使用
Search After或ScrollAPI。
4.3 搜索流程总结图
五、核心概念对比与 Type 演变
5.1 核心概念对比
┌─────────────────────────────────────────────────────────────────────────┐
│ ES 与关系型数据库概念对比 │
├───────────────────┬─────────────────────┬───────────────────────────────┤
│ 关系型数据库 │ Elasticsearch │ 说明 │
├───────────────────┼─────────────────────┼───────────────────────────────┤
│ Database │ Cluster │ 数据库/集群 │
│ Table │ Index │ 表/索引 │
│ Row │ Document │ 行/文档 │
│ Column │ Field │ 列/字段 │
│ Schema │ Mapping │ 表结构/映射 │
│ Index │ Inverted Index │ 索引/倒排索引 │
│ SQL │ Query DSL │ 查询语言 │
└───────────────────┴─────────────────────┴───────────────────────────────┘
5.2 Type 演变历史
- 5.x 及以前:允许一个 Index 下存在多个 Type(类比 Table),但本质上底层字段是扁平化混在一起的,导致数据稀疏(Sparse)问题,影响压缩效率和性能。
- 6.x:强制规定一个 Index 只能有一个 Type,通常默认名为
doc。 - 7.x:Type 概念被彻底废弃(默认为
_doc),API 中 URL 的 type 参数变为可选。 - 8.x:彻底移除 Type 概念。
- 结论:现在设计索引时,严格遵循 “一个 Index 对应一类业务数据” 的原则。
为什么移除?
- 映射爆炸(Mapping Explosion)
- 在使用多类型时,如果不同类型之间有大量不同的字段,这会导致映射的数量急剧增加,进而引发映射爆炸问题。映射爆炸不仅会消耗大量的内存资源,还会降低 Elasticsearch 的性能,尤其是在处理大量数据时
- 字段名冲突
- 在同一个index的不同type中,如果有相同名称但映射类型不同的字段,会造成字段名冲突。这是因为 Elasticsearch 在内部是将这些字段扁平化处理的,而不同类型的相同名称字段可能会导致数据解析和查询时的混乱
- 我们可以和关系型数据库来对比,在同一个数据库中,这些不同的表,可以有名称相同但类型不同的字段。而在 Elasticsearch 同一个index的不同type中,如果有不同document的字段名相同,但是类型不同,就会报错
- 综上所述,Elasticsearch 从 7.x 版本开始废弃类型的主要目的是为了提升系统的性能、避免映射爆炸和字段冲突的问题,以及简化数据模型的设计和管理。这一改变反映了 Elasticsearch 对于提高性能、可维护性和用户体验的持续追求
六、核心功能
6.1 搜索能力矩阵
6.2 聚合分析
七、应用场景
- 日志与监控(ELK Stack):收集服务器日志、应用 Error 日志,快速定位故障(Logstash/Beats + ES + Kibana)。
- 站内搜索:电商商品搜索、论坛帖子搜索、企业知识库检索。
- 大屏可视化/BI:实时统计大盘数据(如双11大屏),利用聚合功能快速出报表。
- 地理位置服务(LBS):查询“附近的酒店”、“方圆5公里内的订单”(Geo-point/Geo-shape)。
八、如何使用 Elasticsearch
8.1 Spring Boot 集成
通常有两种主流方式:
- Spring Data Elasticsearch:封装程度极高,类似 JPA/MyBatis-Plus,通过 Repository 接口操作。
- 优点:开发极快,代码简洁。
- 缺点:灵活性稍差,对复杂 DSL 和版本兼容性控制不如原生客户端细致。
- RestHighLevelClient (官方推荐/传统):基于 HTTP 的原生客户端封装。
- 优点:完全覆盖官方 API,灵活,可控性强。
- 注意:ES 7.15+ 后官方推出了新的
Elasticsearch Java API Client,但RestHighLevelClient依然在存量系统中广泛使用。本文基于此方案进行封装。
<!-- Maven 依赖 -->
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-data-elasticsearch</artifactId>
</dependency>
# application.yml 配置
spring:
elasticsearch:
uris: http://localhost:9200
username: elastic
password: password
此type的作用就是为了兼容6.x、7.x中的type概念,默认是关闭
8.2 原生 API 操作的复杂性
直接使用 RestHighLevelClient 会面临大量样板代码:
- 构建
SearchSourceBuilder、BoolQueryBuilder极其繁琐。 - 需要手动处理
IOException。 - 响应结果解析(Parse)需要从 JSON 层层剥离,非常痛苦。
- 连接管理和配置分散。
// 原生 ES API 操作示例 - 复杂且繁琐
public SearchResponse searchPrograms(String keyword, Integer categoryId) {
// 构建查询条件
BoolQueryBuilder boolQuery = QueryBuilders.boolQuery();
if (StringUtils.isNotBlank(keyword)) {
boolQuery.must(QueryBuilders.matchQuery("title", keyword));
}
if (categoryId != null) {
boolQuery.filter(QueryBuilders.termQuery("categoryId", categoryId));
}
// 构建排序
FieldSortBuilder sortBuilder = SortBuilders.fieldSort("showTime")
.order(SortOrder.ASC);
// 构建搜索源
SearchSourceBuilder sourceBuilder = new SearchSourceBuilder();
sourceBuilder.query(boolQuery);
sourceBuilder.sort(sortBuilder);
sourceBuilder.from(0);
sourceBuilder.size(10);
// 构建搜索请求
SearchRequest searchRequest = new SearchRequest("program-index");
searchRequest.source(sourceBuilder);
// 执行搜索
return restHighLevelClient.search(searchRequest, RequestOptions.DEFAULT);
}
在平时开发中还是使用关系型数据库更加普遍,对于数据库、表、字段的概念更为熟悉,也更加习惯对表概念的操作。
而在操作Elasticsearch时,提供的api其实是很复杂的,各种操作的对象,如SearchSourceBuilder FieldSortBuilder BoolQueryBuilder 等等,操作上其实算不上简单,为了解决这个问题,在springboot操作Elasticsearch的基础上,进一步的封装,使用起来贴近于关系型数据库的方式,操作起来更加的容易上手
九、封装设计
9.1 封装目标
- 统一配置:简化连接参数管理。
- 屏蔽细节:隐藏繁琐的 Builder 构建过程。
- 简化查询:通过 Map 或对象传递参数,自动构建 DSL。
- 结果转换:自动将 ES 的 JSON 结果转为 Java Bean。
- 健壮性:统一异常处理,防止 ES 波动导致服务崩溃。
9.2 封装架构设计
elasticsearch-spring-boot-starter/
├── src/main/java/com/example/elasticsearch/
│ ├── config/
│ │ ├── ElasticsearchProperties.java # 配置属性
│ │ └── ElasticsearchAutoConfiguration.java # 自动配置
│ ├── core/
│ │ ├── ElasticsearchService.java # 核心服务类
│ │ ├── ElasticsearchIndexService.java # 索引操作服务
│ │ └── ElasticsearchDocumentService.java # 文档操作服务
│ ├── query/
│ │ ├── EsQueryBuilder.java # 查询构建器
│ │ ├── EsSearchRequest.java # 搜索请求封装
│ │ └── EsHighlightConfig.java # 高亮配置
│ ├── result/
│ │ ├── PageResult.java # 分页结果
│ │ ├── SearchResult.java # 搜索结果
│ │ └── AggregationResult.java # 聚合结果
│ ├── exception/
│ │ └── ElasticsearchException.java # 自定义异常
│ └── annotation/
│ ├── EsDocument.java # 文档注解
│ └── EsField.java # 字段注解
└── src/main/resources/
└── META-INF/spring.factories
9.3 配置类设计
package com.example.elasticsearch.config;
import lombok.Data;
import org.springframework.boot.context.properties.ConfigurationProperties;
import org.springframework.validation.annotation.Validated;
import javax.validation.constraints.Min;
import javax.validation.constraints.NotEmpty;
import java.util.List;
/**
* Elasticsearch 配置属性类
* 支持集群配置、连接池、认证、重试等
*/
@Data
@Validated
@ConfigurationProperties(prefix = "elasticsearch")
public class ElasticsearchProperties {
/**
* 是否启用 ES
*/
private Boolean enabled = true;
/**
* ES 节点地址列表(支持集群)
*/
@NotEmpty(message = "ES 节点地址不能为空")
private List<String> nodes = List.of("localhost:9200");
/**
* 用户名(可选)
*/
private String username;
/**
* 密码(可选)
*/
private String password;
/**
* 协议:http 或 https
*/
private String scheme = "http";
/**
* 是否启用 Type(兼容 ES 6.x 版本)
*/
private Boolean enableType = false;
/**
* 默认 Type 名称
*/
private String defaultType = "_doc";
/**
* 连接超时时间(毫秒)
*/
@Min(value = 1000, message = "连接超时时间不能小于1000ms")
private Integer connectTimeout = 5000;
/**
* Socket 超时时间(毫秒)
*/
@Min(value = 1000, message = "Socket超时时间不能小于1000ms")
private Integer socketTimeout = 30000;
/**
* 请求超时时间(毫秒)
*/
private Integer connectionRequestTimeout = 5000;
/**
* 最大连接数
*/
@Min(value = 1, message = "最大连接数不能小于1")
private Integer maxConnTotal = 100;
/**
* 每个路由的最大连接数
*/
@Min(value = 1, message = "每个路由最大连接数不能小于1")
private Integer maxConnPerRoute = 50;
/**
* 重试次数
*/
@Min(value = 0, message = "重试次数不能为负数")
private Integer retryTimes = 3;
/**
* 重试间隔(毫秒)
*/
private Long retryInterval = 1000L;
/**
* 是否开启嗅探器
*/
private Boolean enableSniffer = false;
/**
* 嗅探间隔时间(毫秒)
*/
private Long snifferInterval = 60000L;
/**
* 批量操作每批大小
*/
@Min(value = 100, message = "批量操作每批大小不能小于100")
private Integer bulkBatchSize = 1000;
/**
* 批量操作刷新策略:immediate, wait_for, none
*/
private String bulkRefreshPolicy = "none";
/**
* 是否打印 DSL 日志
*/
private Boolean printDsl = false;
/**
* 慢查询阈值(毫秒),超过此值记录警告日志
*/
private Long slowQueryThreshold = 3000L;
}
9.4 自动配置类
package com.example.elasticsearch.config;
import com.example.elasticsearch.core.ElasticsearchDocumentService;
import com.example.elasticsearch.core.ElasticsearchIndexService;
import com.example.elasticsearch.core.ElasticsearchService;
import lombok.extern.slf4j.Slf4j;
import org.apache.http.HttpHost;
import org.apache.http.auth.AuthScope;
import org.apache.http.auth.UsernamePasswordCredentials;
import org.apache.http.client.CredentialsProvider;
import org.apache.http.impl.client.BasicCredentialsProvider;
import org.elasticsearch.client.RestClient;
import org.elasticsearch.client.RestClientBuilder;
import org.elasticsearch.client.RestHighLevelClient;
import org.elasticsearch.client.sniff.Sniffer;
import org.springframework.boot.autoconfigure.condition.ConditionalOnMissingBean;
import org.springframework.boot.autoconfigure.condition.ConditionalOnProperty;
import org.springframework.boot.context.properties.EnableConfigurationProperties;
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import org.springframework.util.StringUtils;
import javax.annotation.PreDestroy;
import java.util.List;
import java.util.stream.Collectors;
/**
* Elasticsearch 自动配置类
*/
@Slf4j
@Configuration
@EnableConfigurationProperties(ElasticsearchProperties.class)
@ConditionalOnProperty(prefix = "elasticsearch", name = "enabled", havingValue = "true", matchIfMissing = true)
public class ElasticsearchAutoConfiguration {
private RestHighLevelClient restHighLevelClient;
private Sniffer sniffer;
@Bean
@ConditionalOnMissingBean
public RestHighLevelClient restHighLevelClient(ElasticsearchProperties properties) {
// 解析节点地址
List<HttpHost> httpHosts = properties.getNodes().stream()
.map(node -> {
String[] parts = node.split(":");
String host = parts[0];
int port = parts.length > 1 ? Integer.parseInt(parts[1]) : 9200;
return new HttpHost(host, port, properties.getScheme());
})
.collect(Collectors.toList());
RestClientBuilder builder = RestClient.builder(
httpHosts.toArray(new HttpHost[0])
);
// 设置请求配置
builder.setRequestConfigCallback(requestConfigBuilder ->
requestConfigBuilder
.setConnectTimeout(properties.getConnectTimeout())
.setSocketTimeout(properties.getSocketTimeout())
.setConnectionRequestTimeout(properties.getConnectionRequestTimeout())
);
// 设置 HTTP 客户端配置
builder.setHttpClientConfigCallback(httpClientBuilder -> {
// 设置连接池
httpClientBuilder.setMaxConnTotal(properties.getMaxConnTotal());
httpClientBuilder.setMaxConnPerRoute(properties.getMaxConnPerRoute());
// 设置认证
if (StringUtils.hasText(properties.getUsername())
&& StringUtils.hasText(properties.getPassword())) {
CredentialsProvider credentialsProvider = new BasicCredentialsProvider();
credentialsProvider.setCredentials(
AuthScope.ANY,
new UsernamePasswordCredentials(
properties.getUsername(),
properties.getPassword()
)
);
httpClientBuilder.setDefaultCredentialsProvider(credentialsProvider);
}
return httpClientBuilder;
});
// 设置失败重试策略
builder.setFailureListener(new RestClient.FailureListener() {
@Override
public void onFailure(org.elasticsearch.client.Node node) {
log.warn("ES 节点 [{}] 连接失败", node.getHost());
}
});
restHighLevelClient = new RestHighLevelClient(builder);
// 启用嗅探器
if (properties.getEnableSniffer()) {
sniffer = Sniffer.builder(restHighLevelClient.getLowLevelClient())
.setSniffIntervalMillis(properties.getSnifferInterval().intValue())
.build();
log.info("ES 嗅探器已启用,间隔: {}ms", properties.getSnifferInterval());
}
log.info("ES 客户端初始化成功, 节点: {}", properties.getNodes());
return restHighLevelClient;
}
@Bean
@ConditionalOnMissingBean
public ElasticsearchService elasticsearchService(
RestHighLevelClient client,
ElasticsearchProperties properties) {
return new ElasticsearchService(client, properties);
}
@Bean
@ConditionalOnMissingBean
public ElasticsearchIndexService elasticsearchIndexService(
RestHighLevelClient client,
ElasticsearchProperties properties) {
return new ElasticsearchIndexService(client, properties);
}
@Bean
@ConditionalOnMissingBean
public ElasticsearchDocumentService elasticsearchDocumentService(
RestHighLevelClient client,
ElasticsearchProperties properties) {
return new ElasticsearchDocumentService(client, properties);
}
@PreDestroy
public void destroy() {
try {
if (sniffer != null) {
sniffer.close();
log.info("ES 嗅探器已关闭");
}
if (restHighLevelClient != null) {
restHighLevelClient.close();
log.info("ES 客户端已关闭");
}
} catch (Exception e) {
log.error("关闭 ES 客户端失败", e);
}
}
}
9.5 自定义异常类
package com.example.elasticsearch.exception;
import lombok.Getter;
/**
* Elasticsearch 自定义异常
*/
@Getter
public class ElasticsearchException extends RuntimeException {
private static final long serialVersionUID = 1L;
/**
* 索引名称
*/
private String index;
/**
* 操作类型
*/
private String operation;
/**
* 错误码
*/
private String errorCode;
/**
* 是否可重试
*/
private boolean retryable;
public ElasticsearchException(String message) {
super(message);
this.retryable = false;
}
public ElasticsearchException(String message, Throwable cause) {
super(message, cause);
this.retryable = isRetryableException(cause);
}
public ElasticsearchException(String operation, String index, String message) {
super(String.format("[%s] 索引 [%s] 操作失败: %s", operation, index, message));
this.operation = operation;
this.index = index;
this.errorCode = operation + "_ERROR";
}
public ElasticsearchException(String operation, String index, Throwable cause) {
super(String.format("[%s] 索引 [%s] 操作失败: %s",
operation, index, cause.getMessage()), cause);
this.operation = operation;
this.index = index;
this.errorCode = operation + "_ERROR";
this.retryable = isRetryableException(cause);
}
/**
* 判断异常是否可重试
*/
private boolean isRetryableException(Throwable cause) {
if (cause == null) {
return false;
}
String message = cause.getMessage();
if (message == null) {
return false;
}
// 可重试的异常类型
return message.contains("Connection refused")
|| message.contains("Connection reset")
|| message.contains("Connection timed out")
|| message.contains("Read timed out")
|| message.contains("No route to host")
|| message.contains("Service Unavailable")
|| message.contains("circuit_breaking_exception");
}
/**
* 创建索引不存在异常
*/
public static ElasticsearchException indexNotFound(String index) {
ElasticsearchException ex = new ElasticsearchException(
"INDEX_NOT_FOUND", index, "索引不存在");
ex.errorCode = "INDEX_NOT_FOUND";
return ex;
}
/**
* 创建文档不存在异常
*/
public static ElasticsearchException documentNotFound(String index, String id) {
ElasticsearchException ex = new ElasticsearchException(
"DOCUMENT_NOT_FOUND", index,
String.format("文档 [%s] 不存在", id));
ex.errorCode = "DOCUMENT_NOT_FOUND";
return ex;
}
/**
* 创建参数校验异常
*/
public static ElasticsearchException invalidParameter(String message) {
ElasticsearchException ex = new ElasticsearchException(message);
ex.errorCode = "INVALID_PARAMETER";
return ex;
}
}
9.6 结果封装类
package com.example.elasticsearch.result;
import lombok.AllArgsConstructor;
import lombok.Data;
import lombok.NoArgsConstructor;
import java.util.Collections;
import java.util.List;
/**
* 分页结果封装
*/
@Data
@NoArgsConstructor
@AllArgsConstructor
public class PageResult<T> {
/**
* 数据列表
*/
private List<T> records;
/**
* 总记录数
*/
private Long total;
/**
* 当前页码
*/
private Integer pageNum;
/**
* 每页大小
*/
private Integer pageSize;
/**
* 总页数
*/
private Integer totalPages;
/**
* 是否有下一页
*/
private Boolean hasNext;
/**
* 是否有上一页
*/
private Boolean hasPrevious;
public PageResult(List<T> records, Long total, Integer pageNum, Integer pageSize) {
this.records = records;
this.total = total;
this.pageNum = pageNum;
this.pageSize = pageSize;
this.totalPages = (int) Math.ceil((double) total / pageSize);
this.hasNext = pageNum < totalPages;
this.hasPrevious = pageNum > 1;
}
/**
* 空结果
*/
public static <T> PageResult<T> empty(Integer pageNum, Integer pageSize) {
return new PageResult<>(Collections.emptyList(), 0L, pageNum, pageSize);
}
}
package com.example.elasticsearch.result;
import lombok.Data;
import java.util.List;
import java.util.Map;
/**
* 搜索结果封装(包含高亮、评分等信息)
*/
@Data
public class SearchResult<T> {
/**
* 文档 ID
*/
private String documentId;
/**
* 数据对象
*/
private T source;
/**
* 评分
*/
private Float score;
/**
* 高亮字段
*/
private Map<String, List<String>> highlight;
/**
* 排序值(用于深度分页)
*/
private Object[] sortValues;
public SearchResult(String documentId, T source) {
this.documentId = documentId;
this.source = source;
}
public SearchResult(String documentId, T source, Float score,
Map<String, List<String>> highlight) {
this.documentId = documentId;
this.source = source;
this.score = score;
this.highlight = highlight;
}
}
package com.example.elasticsearch.result;
import lombok.Data;
import java.util.List;
import java.util.Map;
/**
* 聚合结果封装
*/
@Data
public class AggregationResult {
/**
* 聚合名称
*/
private String name;
/**
* 桶数据(terms 聚合)
*/
private List<BucketData> buckets;
/**
* 数值(sum、avg、max、min 等)
*/
private Double value;
@Data
public static class BucketData {
private String key;
private Long docCount;
private Map<String, Object> subAggregations;
}
}
9.7 查询构建器
package com.example.elasticsearch.query;
import lombok.Data;
import org.elasticsearch.index.query.*;
import org.elasticsearch.search.aggregations.AggregationBuilder;
import org.elasticsearch.search.aggregations.AggregationBuilders;
import org.elasticsearch.search.builder.SearchSourceBuilder;
import org.elasticsearch.search.fetch.subphase.highlight.HighlightBuilder;
import org.elasticsearch.search.sort.SortBuilder;
import org.elasticsearch.search.sort.SortBuilders;
import org.elasticsearch.search.sort.SortOrder;
import org.springframework.util.CollectionUtils;
import org.springframework.util.StringUtils;
import java.util.*;
/**
* ES 查询构建器
* 使用链式调用构建复杂查询
*/
@Data
public class EsQueryBuilder {
/**
* 索引名称
*/
private String indexName;
/**
* Bool 查询条件
*/
private BoolQueryBuilder boolQuery;
/**
* 排序条件
*/
private List<SortBuilder<?>> sorts;
/**
* 高亮配置
*/
private HighlightBuilder highlightBuilder;
/**
* 聚合配置
*/
private List<AggregationBuilder> aggregations;
/**
* 分页参数
*/
private Integer from;
private Integer size;
/**
* 返回字段(包含)
*/
private String[] includes;
/**
* 排除字段
*/
private String[] excludes;
/**
* Search After(深度分页)
*/
private Object[] searchAfter;
/**
* 是否追踪总数
*/
private Boolean trackTotalHits = true;
private EsQueryBuilder() {
this.boolQuery = QueryBuilders.boolQuery();
this.sorts = new ArrayList<>();
this.aggregations = new ArrayList<>();
}
/**
* 创建构建器
*/
public static EsQueryBuilder builder(String indexName) {
EsQueryBuilder builder = new EsQueryBuilder();
builder.indexName = indexName;
return builder;
}
// ==================== Must 条件(必须匹配)====================
/**
* 精确匹配(term)
*/
public EsQueryBuilder term(String field, Object value) {
if (value != null) {
boolQuery.must(QueryBuilders.termQuery(field, value));
}
return this;
}
/**
* 多值匹配(terms)
*/
public EsQueryBuilder terms(String field, Collection<?> values) {
if (!CollectionUtils.isEmpty(values)) {
boolQuery.must(QueryBuilders.termsQuery(field, values));
}
return this;
}
/**
* 全文匹配(match)
*/
public EsQueryBuilder match(String field, Object value) {
if (value != null && StringUtils.hasText(value.toString())) {
boolQuery.must(QueryBuilders.matchQuery(field, value));
}
return this;
}
/**
* 短语匹配(match_phrase)
*/
public EsQueryBuilder matchPhrase(String field, Object value) {
if (value != null && StringUtils.hasText(value.toString())) {
boolQuery.must(QueryBuilders.matchPhraseQuery(field, value));
}
return this;
}
/**
* 多字段匹配(multi_match)
*/
public EsQueryBuilder multiMatch(Object value, String... fields) {
if (value != null && StringUtils.hasText(value.toString()) && fields.length > 0) {
boolQuery.must(QueryBuilders.multiMatchQuery(value, fields));
}
return this;
}
/**
* 前缀匹配(prefix)
*/
public EsQueryBuilder prefix(String field, String prefix) {
if (StringUtils.hasText(prefix)) {
boolQuery.must(QueryBuilders.prefixQuery(field, prefix));
}
return this;
}
/**
* 通配符匹配(wildcard)
*/
public EsQueryBuilder wildcard(String field, String pattern) {
if (StringUtils.hasText(pattern)) {
boolQuery.must(QueryBuilders.wildcardQuery(field, pattern));
}
return this;
}
/**
* 范围查询
*/
public EsQueryBuilder range(String field, Object gte, Object lte) {
RangeQueryBuilder rangeQuery = QueryBuilders.rangeQuery(field);
if (gte != null) {
rangeQuery.gte(gte);
}
if (lte != null) {
rangeQuery.lte(lte);
}
if (gte != null || lte != null) {
boolQuery.must(rangeQuery);
}
return this;
}
/**
* 范围查询(大于)
*/
public EsQueryBuilder gt(String field, Object value) {
if (value != null) {
boolQuery.must(QueryBuilders.rangeQuery(field).gt(value));
}
return this;
}
/**
* 范围查询(大于等于)
*/
public EsQueryBuilder gte(String field, Object value) {
if (value != null) {
boolQuery.must(QueryBuilders.rangeQuery(field).gte(value));
}
return this;
}
/**
* 范围查询(小于)
*/
public EsQueryBuilder lt(String field, Object value) {
if (value != null) {
boolQuery.must(QueryBuilders.rangeQuery(field).lt(value));
}
return this;
}
/**
* 范围查询(小于等于)
*/
public EsQueryBuilder lte(String field, Object value) {
if (value != null) {
boolQuery.must(QueryBuilders.rangeQuery(field).lte(value));
}
return this;
}
/**
* 存在字段查询
*/
public EsQueryBuilder exists(String field) {
boolQuery.must(QueryBuilders.existsQuery(field));
return this;
}
// ==================== Filter 条件(过滤,不计算评分)====================
/**
* Filter - 精确匹配
*/
public EsQueryBuilder filterTerm(String field, Object value) {
if (value != null) {
boolQuery.filter(QueryBuilders.termQuery(field, value));
}
return this;
}
/**
* Filter - 多值匹配
*/
public EsQueryBuilder filterTerms(String field, Collection<?> values) {
if (!CollectionUtils.isEmpty(values)) {
boolQuery.filter(QueryBuilders.termsQuery(field, values));
}
return this;
}
/**
* Filter - 范围查询
*/
public EsQueryBuilder filterRange(String field, Object gte, Object lte) {
RangeQueryBuilder rangeQuery = QueryBuilders.rangeQuery(field);
if (gte != null) {
rangeQuery.gte(gte);
}
if (lte != null) {
rangeQuery.lte(lte);
}
if (gte != null || lte != null) {
boolQuery.filter(rangeQuery);
}
return this;
}
// ==================== Should 条件(或条件)====================
/**
* Should - 至少匹配一个
*/
public EsQueryBuilder should(QueryBuilder... queries) {
for (QueryBuilder query : queries) {
boolQuery.should(query);
}
return this;
}
/**
* Should - 多字段或查询
*/
public EsQueryBuilder shouldMatch(String value, String... fields) {
if (StringUtils.hasText(value) && fields.length > 0) {
for (String field : fields) {
boolQuery.should(QueryBuilders.matchQuery(field, value));
}
boolQuery.minimumShouldMatch(1);
}
return this;
}
/**
* 设置最小 should 匹配数
*/
public EsQueryBuilder minimumShouldMatch(int count) {
boolQuery.minimumShouldMatch(count);
return this;
}
// ==================== MustNot 条件(必须不匹配)====================
/**
* MustNot - 精确匹配
*/
public EsQueryBuilder mustNotTerm(String field, Object value) {
if (value != null) {
boolQuery.mustNot(QueryBuilders.termQuery(field, value));
}
return this;
}
/**
* MustNot - 多值匹配
*/
public EsQueryBuilder mustNotTerms(String field, Collection<?> values) {
if (!CollectionUtils.isEmpty(values)) {
boolQuery.mustNot(QueryBuilders.termsQuery(field, values));
}
return this;
}
// ==================== 嵌套查询 ====================
/**
* 嵌套查询
*/
public EsQueryBuilder nested(String path, QueryBuilder query) {
boolQuery.must(QueryBuilders.nestedQuery(path, query,
org.apache.lucene.search.join.ScoreMode.Avg));
return this;
}
/**
* 添加自定义查询条件
*/
public EsQueryBuilder must(QueryBuilder query) {
boolQuery.must(query);
return this;
}
/**
* 添加自定义过滤条件
*/
public EsQueryBuilder filter(QueryBuilder query) {
boolQuery.filter(query);
return this;
}
// ==================== 排序 ====================
/**
* 添加排序
*/
public EsQueryBuilder sort(String field, SortOrder order) {
sorts.add(SortBuilders.fieldSort(field).order(order));
return this;
}
/**
* 按评分排序
*/
public EsQueryBuilder sortByScore(SortOrder order) {
sorts.add(SortBuilders.scoreSort().order(order));
return this;
}
/**
* 多字段排序
*/
public EsQueryBuilder sorts(Map<String, SortOrder> sortMap) {
sortMap.forEach((field, order) ->
sorts.add(SortBuilders.fieldSort(field).order(order)));
return this;
}
// ==================== 分页 ====================
/**
* 设置分页
*/
public EsQueryBuilder page(int pageNum, int pageSize) {
this.from = (pageNum - 1) * pageSize;
this.size = pageSize;
return this;
}
/**
* 设置起始位置和大小
*/
public EsQueryBuilder fromSize(int from, int size) {
this.from = from;
this.size = size;
return this;
}
/**
* 深度分页(Search After)
*/
public EsQueryBuilder searchAfter(Object[] values) {
this.searchAfter = values;
return this;
}
// ==================== 高亮 ====================
/**
* 添加高亮字段
*/
public EsQueryBuilder highlight(String... fields) {
if (fields.length > 0) {
highlightBuilder = new HighlightBuilder();
for (String field : fields) {
highlightBuilder.field(field);
}
highlightBuilder.preTags("<em class='highlight'>");
highlightBuilder.postTags("</em>");
}
return this;
}
/**
* 自定义高亮配置
*/
public EsQueryBuilder highlight(String preTag, String postTag, String... fields) {
if (fields.length > 0) {
highlightBuilder = new HighlightBuilder();
for (String field : fields) {
highlightBuilder.field(field);
}
highlightBuilder.preTags(preTag);
highlightBuilder.postTags(postTag);
}
return this;
}
// ==================== 聚合 ====================
/**
* Terms 聚合
*/
public EsQueryBuilder termsAggregation(String name, String field, int size) {
aggregations.add(AggregationBuilders.terms(name).field(field).size(size));
return this;
}
/**
* Sum 聚合
*/
public EsQueryBuilder sumAggregation(String name, String field) {
aggregations.add(AggregationBuilders.sum(name).field(field));
return this;
}
/**
* Avg 聚合
*/
public EsQueryBuilder avgAggregation(String name, String field) {
aggregations.add(AggregationBuilders.avg(name).field(field));
return this;
}
/**
* Max 聚合
*/
public EsQueryBuilder maxAggregation(String name, String field) {
aggregations.add(AggregationBuilders.max(name).field(field));
return this;
}
/**
* Min 聚合
*/
public EsQueryBuilder minAggregation(String name, String field) {
aggregations.add(AggregationBuilders.min(name).field(field));
return this;
}
/**
* 日期直方图聚合
*/
public EsQueryBuilder dateHistogramAggregation(String name, String field,
String interval) {
aggregations.add(AggregationBuilders.dateHistogram(name)
.field(field)
.calendarInterval(new org.elasticsearch.search.aggregations.bucket
.histogram.DateHistogramInterval(interval)));
return this;
}
/**
* 添加自定义聚合
*/
public EsQueryBuilder aggregation(AggregationBuilder aggregation) {
aggregations.add(aggregation);
return this;
}
// ==================== 返回字段 ====================
/**
* 指定返回字段
*/
public EsQueryBuilder includes(String... fields) {
this.includes = fields;
return this;
}
/**
* 排除字段
*/
public EsQueryBuilder excludes(String... fields) {
this.excludes = fields;
return this;
}
/**
* 是否追踪总数
*/
public EsQueryBuilder trackTotalHits(boolean track) {
this.trackTotalHits = track;
return this;
}
// ==================== 构建 ====================
/**
* 构建 SearchSourceBuilder
*/
public SearchSourceBuilder build() {
SearchSourceBuilder sourceBuilder = new SearchSourceBuilder();
// 查询条件
sourceBuilder.query(boolQuery);
// 分页
if (from != null) {
sourceBuilder.from(from);
}
if (size != null) {
sourceBuilder.size(size);
}
// 排序
for (SortBuilder<?> sort : sorts) {
sourceBuilder.sort(sort);
}
// 高亮
if (highlightBuilder != null) {
sourceBuilder.highlighter(highlightBuilder);
}
// 聚合
for (AggregationBuilder aggregation : aggregations) {
sourceBuilder.aggregation(aggregation);
}
// 返回字段
if (includes != null || excludes != null) {
sourceBuilder.fetchSource(includes, excludes);
}
// Search After
if (searchAfter != null) {
sourceBuilder.searchAfter(searchAfter);
}
// 追踪总数
sourceBuilder.trackTotalHits(trackTotalHits);
return sourceBuilder;
}
}
9.8 重试工具类
package com.example.elasticsearch.util;
import com.example.elasticsearch.exception.ElasticsearchException;
import lombok.extern.slf4j.Slf4j;
import java.util.concurrent.Callable;
import java.util.function.Predicate;
/**
* 重试工具类
*/
@Slf4j
public class RetryUtil {
/**
* 执行带重试的操作
*
* @param callable 要执行的操作
* @param maxRetries 最大重试次数
* @param retryInterval 重试间隔(毫秒)
* @param retryOn 判断是否需要重试的条件
* @param operationName 操作名称(用于日志)
* @return 操作结果
*/
public static <T> T executeWithRetry(
Callable<T> callable,
int maxRetries,
long retryInterval,
Predicate<Exception> retryOn,
String operationName) {
Exception lastException = null;
for (int attempt = 0; attempt <= maxRetries; attempt++) {
try {
return callable.call();
} catch (Exception e) {
lastException = e;
// 判断是否需要重试
if (attempt < maxRetries && retryOn.test(e)) {
log.warn("[{}] 操作失败,第 {}/{} 次重试,错误: {}",
operationName, attempt + 1, maxRetries, e.getMessage());
try {
Thread.sleep(retryInterval * (attempt + 1)); // 指数退避
} catch (InterruptedException ie) {
Thread.currentThread().interrupt();
throw new ElasticsearchException("重试被中断", ie);
}
} else {
break;
}
}
}
log.error("[{}] 操作失败,已达最大重试次数", operationName, lastException);
if (lastException instanceof ElasticsearchException) {
throw (ElasticsearchException) lastException;
}
throw new ElasticsearchException(operationName + " 操作失败", lastException);
}
/**
* 执行带重试的操作(无返回值)
*/
public static void executeWithRetry(
Runnable runnable,
int maxRetries,
long retryInterval,
Predicate<Exception> retryOn,
String operationName) {
executeWithRetry(() -> {
runnable.run();
return null;
}, maxRetries, retryInterval, retryOn, operationName);
}
/**
* 默认的重试条件判断
*/
public static Predicate<Exception> defaultRetryCondition() {
return e -> {
if (e instanceof ElasticsearchException) {
return ((ElasticsearchException) e).isRetryable();
}
String message = e.getMessage();
if (message == null) {
return false;
}
return message.contains("Connection")
|| message.contains("timed out")
|| message.contains("Unavailable");
};
}
}
9.9 参数校验工具类
package com.example.elasticsearch.util;
import com.example.elasticsearch.exception.ElasticsearchException;
import org.springframework.util.CollectionUtils;
import org.springframework.util.StringUtils;
import java.util.Collection;
/**
* 参数校验工具类
*/
public class ParamValidator {
private ParamValidator() {}
/**
* 校验索引名称
*/
public static void validateIndexName(String indexName) {
if (!StringUtils.hasText(indexName)) {
throw ElasticsearchException.invalidParameter("索引名称不能为空");
}
if (indexName.contains(" ")) {
throw ElasticsearchException.invalidParameter("索引名称不能包含空格");
}
if (!indexName.equals(indexName.toLowerCase())) {
throw ElasticsearchException.invalidParameter("索引名称必须小写");
}
}
/**
* 校验文档ID
*/
public static void validateDocumentId(String documentId) {
if (!StringUtils.hasText(documentId)) {
throw ElasticsearchException.invalidParameter("文档ID不能为空");
}
}
/**
* 校验分页参数
*/
public static void validatePageParam(int pageNum, int pageSize) {
if (pageNum < 1) {
throw ElasticsearchException.invalidParameter("页码必须大于0");
}
if (pageSize < 1 || pageSize > 10000) {
throw ElasticsearchException.invalidParameter("每页大小必须在1-10000之间");
}
// ES 默认限制 from + size <= 10000
if ((long) (pageNum - 1) * pageSize + pageSize > 10000) {
throw ElasticsearchException.invalidParameter(
"分页深度超出限制,请使用 searchAfter 方式");
}
}
/**
* 校验批量数据
*/
public static void validateBatchData(Collection<?> dataList) {
if (CollectionUtils.isEmpty(dataList)) {
throw ElasticsearchException.invalidParameter("批量数据不能为空");
}
}
/**
* 校验非空
*/
public static void notNull(Object object, String message) {
if (object == null) {
throw ElasticsearchException.invalidParameter(message);
}
}
/**
* 校验字符串非空
*/
public static void notBlank(String str, String message) {
if (!StringUtils.hasText(str)) {
throw ElasticsearchException.invalidParameter(message);
}
}
}
9.10 核心 Service 封装
package com.example.elasticsearch.core;
import com.alibaba.fastjson.JSON;
import com.example.elasticsearch.config.ElasticsearchProperties;
import com.example.elasticsearch.exception.ElasticsearchException;
import com.example.elasticsearch.query.EsQueryBuilder;
import com.example.elasticsearch.result.AggregationResult;
import com.example.elasticsearch.result.PageResult;
import com.example.elasticsearch.result.ScrollResult;
import com.example.elasticsearch.result.SearchResult;
import com.example.elasticsearch.util.ParamValidator;
import com.example.elasticsearch.util.RetryUtil;
import lombok.extern.slf4j.Slf4j;
import org.elasticsearch.action.search.*;
import org.elasticsearch.client.RequestOptions;
import org.elasticsearch.client.RestHighLevelClient;
import org.elasticsearch.common.text.Text;
import org.elasticsearch.common.unit.TimeValue;
import org.elasticsearch.index.query.QueryBuilder;
import org.elasticsearch.search.SearchHit;
import org.elasticsearch.search.SearchHits;
import org.elasticsearch.search.aggregations.Aggregation;
import org.elasticsearch.search.aggregations.bucket.terms.Terms;
import org.elasticsearch.search.aggregations.metrics.*;
import org.elasticsearch.search.builder.SearchSourceBuilder;
import org.elasticsearch.search.fetch.subphase.highlight.HighlightField;
import java.io.IOException;
import java.util.*;
import java.util.stream.Collectors;
/**
* Elasticsearch 核心搜索服务
*
* 特性:
* - 完善的参数校验
* - 自动重试机制
* - 慢查询日志
* - 滚动查询支持
* - 空值安全处理
*/
@Slf4j
public class ElasticsearchService {
private final RestHighLevelClient client;
private final ElasticsearchProperties properties;
public ElasticsearchService(RestHighLevelClient client,
ElasticsearchProperties properties) {
this.client = client;
this.properties = properties;
}
// ==================== 基础查询 ====================
/**
* 简单查询 - 根据单个字段精确匹配
*
* @param indexName 索引名称
* @param field 字段名
* @param value 字段值
* @param clazz 返回类型
* @return 匹配的文档列表
*/
public <T> List<T> query(String indexName, String field, Object value,
Class<T> clazz) {
ParamValidator.validateIndexName(indexName);
ParamValidator.notBlank(field, "查询字段不能为空");
EsQueryBuilder queryBuilder = EsQueryBuilder.builder(indexName)
.term(field, value);
return search(queryBuilder, clazz);
}
/**
* 多条件查询 - 根据多个字段匹配
*
* @param indexName 索引名称
* @param params 查询参数 (字段名 -> 字段值)
* @param clazz 返回类型
* @return 匹配的文档列表
*/
public <T> List<T> query(String indexName, Map<String, Object> params,
Class<T> clazz) {
ParamValidator.validateIndexName(indexName);
EsQueryBuilder queryBuilder = EsQueryBuilder.builder(indexName);
if (params != null && !params.isEmpty()) {
params.forEach((field, value) -> {
if (value != null) {
if (value instanceof String && ((String) value).length() > 0) {
queryBuilder.match(field, value);
} else if (!(value instanceof String)) {
queryBuilder.filterTerm(field, value);
}
}
});
}
return search(queryBuilder, clazz);
}
/**
* 分页查询
*
* @param indexName 索引名称
* @param params 查询参数
* @param pageNum 页码(从1开始)
* @param pageSize 每页大小
* @param clazz 返回类型
* @return 分页结果
*/
public <T> PageResult<T> queryPage(String indexName, Map<String, Object> params,
int pageNum, int pageSize, Class<T> clazz) {
ParamValidator.validateIndexName(indexName);
ParamValidator.validatePageParam(pageNum, pageSize);
EsQueryBuilder queryBuilder = EsQueryBuilder.builder(indexName)
.page(pageNum, pageSize);
if (params != null && !params.isEmpty()) {
params.forEach((field, value) -> {
if (value != null) {
if (value instanceof String && ((String) value).length() > 0) {
queryBuilder.match(field, value);
} else if (!(value instanceof String)) {
queryBuilder.filterTerm(field, value);
}
}
});
}
return searchPage(queryBuilder, pageNum, pageSize, clazz);
}
// ==================== 高级查询 ====================
/**
* 使用查询构建器执行查询
*
* @param queryBuilder 查询构建器
* @param clazz 返回类型
* @return 匹配的文档列表
*/
public <T> List<T> search(EsQueryBuilder queryBuilder, Class<T> clazz) {
ParamValidator.validateIndexName(queryBuilder.getIndexName());
ParamValidator.notNull(clazz, "返回类型不能为空");
return RetryUtil.executeWithRetry(
() -> doSearch(queryBuilder, clazz),
properties.getRetryTimes(),
properties.getRetryInterval(),
RetryUtil.defaultRetryCondition(),
"SEARCH"
);
}
private <T> List<T> doSearch(EsQueryBuilder queryBuilder, Class<T> clazz)
throws IOException {
long startTime = System.currentTimeMillis();
try {
SearchSourceBuilder sourceBuilder = queryBuilder.build();
SearchRequest searchRequest = new SearchRequest(queryBuilder.getIndexName());
searchRequest.source(sourceBuilder);
// 打印 DSL
if (properties.getPrintDsl()) {
log.info("ES 查询 DSL: {}", sourceBuilder.toString());
}
SearchResponse response = client.search(searchRequest, RequestOptions.DEFAULT);
// 检查响应状态
checkResponseStatus(response, queryBuilder.getIndexName());
return parseHits(response.getHits(), clazz);
} finally {
logSlowQuery(startTime, "SEARCH", queryBuilder.getIndexName());
}
}
/**
* 分页查询
*/
public <T> PageResult<T> searchPage(EsQueryBuilder queryBuilder,
int pageNum, int pageSize, Class<T> clazz) {
ParamValidator.validateIndexName(queryBuilder.getIndexName());
ParamValidator.validatePageParam(pageNum, pageSize);
ParamValidator.notNull(clazz, "返回类型不能为空");
return RetryUtil.executeWithRetry(
() -> doSearchPage(queryBuilder, pageNum, pageSize, clazz),
properties.getRetryTimes(),
properties.getRetryInterval(),
RetryUtil.defaultRetryCondition(),
"SEARCH_PAGE"
);
}
private <T> PageResult<T> doSearchPage(EsQueryBuilder queryBuilder,
int pageNum, int pageSize,
Class<T> clazz) throws IOException {
long startTime = System.currentTimeMillis();
try {
queryBuilder.page(pageNum, pageSize);
SearchSourceBuilder sourceBuilder = queryBuilder.build();
SearchRequest searchRequest = new SearchRequest(queryBuilder.getIndexName());
searchRequest.source(sourceBuilder);
if (properties.getPrintDsl()) {
log.info("ES 分页查询 DSL: {}", sourceBuilder.toString());
}
SearchResponse response = client.search(searchRequest, RequestOptions.DEFAULT);
checkResponseStatus(response, queryBuilder.getIndexName());
List<T> records = parseHits(response.getHits(), clazz);
long total = getTotalHits(response.getHits());
return new PageResult<>(records, total, pageNum, pageSize);
} finally {
logSlowQuery(startTime, "SEARCH_PAGE", queryBuilder.getIndexName());
}
}
/**
* 带高亮的查询
*/
public <T> List<SearchResult<T>> searchWithHighlight(EsQueryBuilder queryBuilder,
Class<T> clazz) {
ParamValidator.validateIndexName(queryBuilder.getIndexName());
return RetryUtil.executeWithRetry(
() -> doSearchWithHighlight(queryBuilder, clazz),
properties.getRetryTimes(),
properties.getRetryInterval(),
RetryUtil.defaultRetryCondition(),
"SEARCH_HIGHLIGHT"
);
}
private <T> List<SearchResult<T>> doSearchWithHighlight(
EsQueryBuilder queryBuilder, Class<T> clazz) throws IOException {
long startTime = System.currentTimeMillis();
try {
SearchSourceBuilder sourceBuilder = queryBuilder.build();
SearchRequest searchRequest = new SearchRequest(queryBuilder.getIndexName());
searchRequest.source(sourceBuilder);
if (properties.getPrintDsl()) {
log.info("ES 高亮查询 DSL: {}", sourceBuilder.toString());
}
SearchResponse response = client.search(searchRequest, RequestOptions.DEFAULT);
checkResponseStatus(response, queryBuilder.getIndexName());
return parseHitsWithHighlight(response.getHits(), clazz);
} finally {
logSlowQuery(startTime, "SEARCH_HIGHLIGHT", queryBuilder.getIndexName());
}
}
/**
* 带高亮的分页查询
*/
public <T> PageResult<SearchResult<T>> searchPageWithHighlight(
EsQueryBuilder queryBuilder, int pageNum, int pageSize, Class<T> clazz) {
ParamValidator.validateIndexName(queryBuilder.getIndexName());
ParamValidator.validatePageParam(pageNum, pageSize);
return RetryUtil.executeWithRetry(
() -> doSearchPageWithHighlight(queryBuilder, pageNum, pageSize, clazz),
properties.getRetryTimes(),
properties.getRetryInterval(),
RetryUtil.defaultRetryCondition(),
"SEARCH_PAGE_HIGHLIGHT"
);
}
private <T> PageResult<SearchResult<T>> doSearchPageWithHighlight(
EsQueryBuilder queryBuilder, int pageNum, int pageSize,
Class<T> clazz) throws IOException {
long startTime = System.currentTimeMillis();
try {
queryBuilder.page(pageNum, pageSize);
SearchSourceBuilder sourceBuilder = queryBuilder.build();
SearchRequest searchRequest = new SearchRequest(queryBuilder.getIndexName());
searchRequest.source(sourceBuilder);
SearchResponse response = client.search(searchRequest, RequestOptions.DEFAULT);
checkResponseStatus(response, queryBuilder.getIndexName());
List<SearchResult<T>> records = parseHitsWithHighlight(response.getHits(), clazz);
long total = getTotalHits(response.getHits());
return new PageResult<>(records, total, pageNum, pageSize);
} finally {
logSlowQuery(startTime, "SEARCH_PAGE_HIGHLIGHT", queryBuilder.getIndexName());
}
}
// ==================== 滚动查询(大数据量) ====================
/**
* 滚动查询 - 初始化
* 适用于导出大量数据的场景
*
* @param queryBuilder 查询构建器
* @param scrollTime 滚动上下文保持时间(分钟)
* @param size 每批大小
* @param clazz 返回类型
* @return 滚动结果(包含scrollId和首批数据)
*/
public <T> ScrollResult<T> scrollSearch(EsQueryBuilder queryBuilder,
int scrollTime, int size,
Class<T> clazz) {
ParamValidator.validateIndexName(queryBuilder.getIndexName());
try {
queryBuilder.fromSize(0, size);
SearchSourceBuilder sourceBuilder = queryBuilder.build();
SearchRequest searchRequest = new SearchRequest(queryBuilder.getIndexName());
searchRequest.source(sourceBuilder);
searchRequest.scroll(TimeValue.timeValueMinutes(scrollTime));
if (properties.getPrintDsl()) {
log.info("ES 滚动查询初始化 DSL: {}", sourceBuilder.toString());
}
SearchResponse response = client.search(searchRequest, RequestOptions.DEFAULT);
checkResponseStatus(response, queryBuilder.getIndexName());
List<T> records = parseHits(response.getHits(), clazz);
long total = getTotalHits(response.getHits());
String scrollId = response.getScrollId();
return new ScrollResult<>(scrollId, records, total, records.size() < size);
} catch (IOException e) {
throw new ElasticsearchException("SCROLL_SEARCH",
queryBuilder.getIndexName(), e);
}
}
/**
* 滚动查询 - 继续获取下一批
*
* @param scrollId 滚动ID
* @param scrollTime 滚动上下文保持时间(分钟)
* @param size 每批大小
* @param clazz 返回类型
* @return 滚动结果
*/
public <T> ScrollResult<T> scrollNext(String scrollId, int scrollTime,
int size, Class<T> clazz) {
ParamValidator.notBlank(scrollId, "scrollId不能为空");
try {
SearchScrollRequest scrollRequest = new SearchScrollRequest(scrollId);
scrollRequest.scroll(TimeValue.timeValueMinutes(scrollTime));
SearchResponse response = client.scroll(scrollRequest, RequestOptions.DEFAULT);
List<T> records = parseHits(response.getHits(), clazz);
String newScrollId = response.getScrollId();
return new ScrollResult<>(newScrollId, records,
getTotalHits(response.getHits()), records.size() < size);
} catch (IOException e) {
throw new ElasticsearchException("滚动查询失败", e);
}
}
/**
* 清除滚动上下文
*
* @param scrollIds 滚动ID列表
*/
public void clearScroll(String... scrollIds) {
if (scrollIds == null || scrollIds.length == 0) {
return;
}
try {
ClearScrollRequest clearScrollRequest = new ClearScrollRequest();
clearScrollRequest.scrollIds(Arrays.asList(scrollIds));
ClearScrollResponse response = client.clearScroll(
clearScrollRequest, RequestOptions.DEFAULT);
if (!response.isSucceeded()) {
log.warn("清除滚动上下文失败");
}
} catch (IOException e) {
log.warn("清除滚动上下文异常", e);
}
}
// ==================== 聚合查询 ====================
/**
* 聚合查询
*/
public Map<String, AggregationResult> searchAggregation(EsQueryBuilder queryBuilder) {
ParamValidator.validateIndexName(queryBuilder.getIndexName());
return RetryUtil.executeWithRetry(
() -> doSearchAggregation(queryBuilder),
properties.getRetryTimes(),
properties.getRetryInterval(),
RetryUtil.defaultRetryCondition(),
"SEARCH_AGGREGATION"
);
}
private Map<String, AggregationResult> doSearchAggregation(
EsQueryBuilder queryBuilder) throws IOException {
long startTime = System.currentTimeMillis();
try {
queryBuilder.fromSize(0, 0);
SearchSourceBuilder sourceBuilder = queryBuilder.build();
SearchRequest searchRequest = new SearchRequest(queryBuilder.getIndexName());
searchRequest.source(sourceBuilder);
if (properties.getPrintDsl()) {
log.info("ES 聚合查询 DSL: {}", sourceBuilder.toString());
}
SearchResponse response = client.search(searchRequest, RequestOptions.DEFAULT);
checkResponseStatus(response, queryBuilder.getIndexName());
return parseAggregations(response.getAggregations());
} finally {
logSlowQuery(startTime, "SEARCH_AGGREGATION", queryBuilder.getIndexName());
}
}
// ==================== Search After(深度分页) ====================
/**
* Search After 查询
* 适用于深度分页场景,避免 from+size 的 10000 限制
*
* @param queryBuilder 查询构建器
* @param searchAfterValues 上一页最后一条记录的排序值
* @param size 每页大小
* @param clazz 返回类型
* @return 搜索结果列表
*/
public <T> List<SearchResult<T>> searchAfter(EsQueryBuilder queryBuilder,
Object[] searchAfterValues,
int size, Class<T> clazz) {
ParamValidator.validateIndexName(queryBuilder.getIndexName());
return RetryUtil.executeWithRetry(
() -> doSearchAfter(queryBuilder, searchAfterValues, size, clazz),
properties.getRetryTimes(),
properties.getRetryInterval(),
RetryUtil.defaultRetryCondition(),
"SEARCH_AFTER"
);
}
private <T> List<SearchResult<T>> doSearchAfter(
EsQueryBuilder queryBuilder, Object[] searchAfterValues,
int size, Class<T> clazz) throws IOException {
long startTime = System.currentTimeMillis();
try {
queryBuilder.fromSize(0, size);
if (searchAfterValues != null && searchAfterValues.length > 0) {
queryBuilder.searchAfter(searchAfterValues);
}
SearchSourceBuilder sourceBuilder = queryBuilder.build();
SearchRequest searchRequest = new SearchRequest(queryBuilder.getIndexName());
searchRequest.source(sourceBuilder);
if (properties.getPrintDsl()) {
log.info("ES Search After DSL: {}", sourceBuilder.toString());
}
SearchResponse response = client.search(searchRequest, RequestOptions.DEFAULT);
checkResponseStatus(response, queryBuilder.getIndexName());
return parseHitsWithHighlight(response.getHits(), clazz);
} finally {
logSlowQuery(startTime, "SEARCH_AFTER", queryBuilder.getIndexName());
}
}
// ==================== 统计查询 ====================
/**
* 统计符合条件的文档数量
*/
public long count(EsQueryBuilder queryBuilder) {
ParamValidator.validateIndexName(queryBuilder.getIndexName());
return RetryUtil.executeWithRetry(
() -> doCount(queryBuilder),
properties.getRetryTimes(),
properties.getRetryInterval(),
RetryUtil.defaultRetryCondition(),
"COUNT"
);
}
private long doCount(EsQueryBuilder queryBuilder) throws IOException {
queryBuilder.fromSize(0, 0);
SearchSourceBuilder sourceBuilder = queryBuilder.build();
SearchRequest searchRequest = new SearchRequest(queryBuilder.getIndexName());
searchRequest.source(sourceBuilder);
SearchResponse response = client.search(searchRequest, RequestOptions.DEFAULT);
checkResponseStatus(response, queryBuilder.getIndexName());
return getTotalHits(response.getHits());
}
/**
* 判断是否存在符合条件的文档
*/
public boolean exists(EsQueryBuilder queryBuilder) {
return count(queryBuilder) > 0;
}
// ==================== 原生查询 ====================
/**
* 执行原生 QueryBuilder 查询
* 适用于封装方法无法满足的复杂场景
*/
public <T> List<T> searchByQueryBuilder(String indexName,
QueryBuilder queryBuilder,
int size, Class<T> clazz) {
ParamValidator.validateIndexName(indexName);
ParamValidator.notNull(queryBuilder, "查询条件不能为空");
try {
SearchSourceBuilder sourceBuilder = new SearchSourceBuilder();
sourceBuilder.query(queryBuilder);
sourceBuilder.size(size);
sourceBuilder.trackTotalHits(true);
SearchRequest searchRequest = new SearchRequest(indexName);
searchRequest.source(sourceBuilder);
SearchResponse response = client.search(searchRequest, RequestOptions.DEFAULT);
checkResponseStatus(response, indexName);
return parseHits(response.getHits(), clazz);
} catch (IOException e) {
throw new ElasticsearchException("SEARCH_BY_QUERY_BUILDER", indexName, e);
}
}
// ==================== 内部工具方法 ====================
/**
* 解析命中结果
*/
private <T> List<T> parseHits(SearchHits hits, Class<T> clazz) {
if (hits == null || hits.getHits() == null) {
return Collections.emptyList();
}
List<T> result = new ArrayList<>();
for (SearchHit hit : hits.getHits()) {
try {
String source = hit.getSourceAsString();
if (source != null && !source.isEmpty()) {
T obj = JSON.parseObject(source, clazz);
result.add(obj);
}
} catch (Exception e) {
log.warn("解析文档失败, id: {}, error: {}", hit.getId(), e.getMessage());
}
}
return result;
}
/**
* 解析带高亮的命中结果
*/
private <T> List<SearchResult<T>> parseHitsWithHighlight(SearchHits hits,
Class<T> clazz) {
if (hits == null || hits.getHits() == null) {
return Collections.emptyList();
}
List<SearchResult<T>> result = new ArrayList<>();
for (SearchHit hit : hits.getHits()) {
try {
String source = hit.getSourceAsString();
if (source == null || source.isEmpty()) {
continue;
}
T obj = JSON.parseObject(source, clazz);
// 解析高亮
Map<String, List<String>> highlightMap = new HashMap<>();
Map<String, HighlightField> highlightFields = hit.getHighlightFields();
if (highlightFields != null && !highlightFields.isEmpty()) {
highlightFields.forEach((field, highlightField) -> {
if (highlightField.getFragments() != null) {
List<String> fragments = Arrays.stream(highlightField.getFragments())
.map(Text::string)
.collect(Collectors.toList());
highlightMap.put(field, fragments);
}
});
}
SearchResult<T> searchResult = new SearchResult<>(
hit.getId(), obj, hit.getScore(), highlightMap);
searchResult.setSortValues(hit.getSortValues());
result.add(searchResult);
} catch (Exception e) {
log.warn("解析文档失败, id: {}, error: {}", hit.getId(), e.getMessage());
}
}
return result;
}
/**
* 解析聚合结果
*/
private Map<String, AggregationResult> parseAggregations(
org.elasticsearch.search.aggregations.Aggregations aggregations) {
Map<String, AggregationResult> result = new HashMap<>();
if (aggregations == null) {
return result;
}
for (Aggregation aggregation : aggregations) {
try {
AggregationResult aggResult = new AggregationResult();
aggResult.setName(aggregation.getName());
if (aggregation instanceof Terms) {
Terms terms = (Terms) aggregation;
List<AggregationResult.BucketData> buckets = new ArrayList<>();
for (Terms.Bucket bucket : terms.getBuckets()) {
AggregationResult.BucketData bucketData =
new AggregationResult.BucketData();
bucketData.setKey(bucket.getKeyAsString());
bucketData.setDocCount(bucket.getDocCount());
buckets.add(bucketData);
}
aggResult.setBuckets(buckets);
} else if (aggregation instanceof Sum) {
aggResult.setValue(((Sum) aggregation).getValue());
} else if (aggregation instanceof Avg) {
aggResult.setValue(((Avg) aggregation).getValue());
} else if (aggregation instanceof Max) {
aggResult.setValue(((Max) aggregation).getValue());
} else if (aggregation instanceof Min) {
aggResult.setValue(((Min) aggregation).getValue());
} else if (aggregation instanceof ValueCount) {
aggResult.setValue((double) ((ValueCount) aggregation).getValue());
} else if (aggregation instanceof Cardinality) {
aggResult.setValue((double) ((Cardinality) aggregation).getValue());
}
result.put(aggregation.getName(), aggResult);
} catch (Exception e) {
log.warn("解析聚合结果失败, name: {}, error: {}",
aggregation.getName(), e.getMessage());
}
}
return result;
}
/**
* 获取总命中数(兼容不同版本)
*/
private long getTotalHits(SearchHits hits) {
if (hits == null || hits.getTotalHits() == null) {
return 0L;
}
return hits.getTotalHits().value;
}
/**
* 检查响应状态
*/
private void checkResponseStatus(SearchResponse response, String indexName) {
if (response == null) {
throw new ElasticsearchException("SEARCH", indexName, "响应为空");
}
if (response.isTimedOut()) {
log.warn("ES 查询超时, index: {}", indexName);
}
if (response.getShardFailures() != null && response.getShardFailures().length > 0) {
log.warn("ES 查询部分分片失败, index: {}, failures: {}",
indexName, response.getShardFailures().length);
}
}
/**
* 记录慢查询日志
*/
private void logSlowQuery(long startTime, String operation, String indexName) {
long elapsed = System.currentTimeMillis() - startTime;
if (elapsed >= properties.getSlowQueryThreshold()) {
log.warn("[慢查询] 操作: {}, 索引: {}, 耗时: {}ms",
operation, indexName, elapsed);
} else {
log.debug("ES 操作完成, 操作: {}, 索引: {}, 耗时: {}ms",
operation, indexName, elapsed);
}
}
}
9.11 文档操作服务
package com.example.elasticsearch.core;
import com.alibaba.fastjson.JSON;
import com.example.elasticsearch.config.ElasticsearchProperties;
import com.example.elasticsearch.exception.ElasticsearchException;
import com.example.elasticsearch.util.ParamValidator;
import com.example.elasticsearch.util.RetryUtil;
import lombok.extern.slf4j.Slf4j;
import org.elasticsearch.action.DocWriteResponse;
import org.elasticsearch.action.bulk.BulkItemResponse;
import org.elasticsearch.action.bulk.BulkRequest;
import org.elasticsearch.action.bulk.BulkResponse;
import org.elasticsearch.action.delete.DeleteRequest;
import org.elasticsearch.action.delete.DeleteResponse;
import org.elasticsearch.action.get.*;
import org.elasticsearch.action.index.IndexRequest;
import org.elasticsearch.action.index.IndexResponse;
import org.elasticsearch.action.support.WriteRequest;
import org.elasticsearch.action.update.UpdateRequest;
import org.elasticsearch.action.update.UpdateResponse;
import org.elasticsearch.client.RequestOptions;
import org.elasticsearch.client.RestHighLevelClient;
import org.elasticsearch.common.xcontent.XContentType;
import org.elasticsearch.index.query.QueryBuilder;
import org.elasticsearch.index.reindex.BulkByScrollResponse;
import org.elasticsearch.index.reindex.DeleteByQueryRequest;
import org.elasticsearch.index.reindex.UpdateByQueryRequest;
import org.elasticsearch.script.Script;
import org.elasticsearch.script.ScriptType;
import org.springframework.util.CollectionUtils;
import org.springframework.util.StringUtils;
import java.io.IOException;
import java.util.*;
/**
* Elasticsearch 文档操作服务(增强版)
*
* 特性:
* - 完善的参数校验
* - 批量操作分批处理
* - 自动重试机制
* - 详细的操作日志
*/
@Slf4j
public class ElasticsearchDocumentService {
private final RestHighLevelClient client;
private final ElasticsearchProperties properties;
public ElasticsearchDocumentService(RestHighLevelClient client,
ElasticsearchProperties properties) {
this.client = client;
this.properties = properties;
}
// ==================== 新增操作 ====================
/**
* 添加文档(自动生成ID)
*/
public String add(String indexName, Object data) {
return add(indexName, null, data);
}
/**
* 添加文档(指定ID)
*/
public String add(String indexName, String id, Object data) {
ParamValidator.validateIndexName(indexName);
ParamValidator.notNull(data, "文档数据不能为空");
return RetryUtil.executeWithRetry(
() -> doAdd(indexName, id, data),
properties.getRetryTimes(),
properties.getRetryInterval(),
RetryUtil.defaultRetryCondition(),
"ADD_DOCUMENT"
);
}
private String doAdd(String indexName, String id, Object data) throws IOException {
IndexRequest request = new IndexRequest(indexName);
if (StringUtils.hasText(id)) {
request.id(id);
}
String jsonData = JSON.toJSONString(data);
request.source(jsonData, XContentType.JSON);
// 兼容 ES 6.x
if (properties.getEnableType()) {
request.type(properties.getDefaultType());
}
// 设置刷新策略
setRefreshPolicy(request);
IndexResponse response = client.index(request, RequestOptions.DEFAULT);
log.info("添加文档成功, index: {}, id: {}", indexName, response.getId());
return response.getId();
}
/**
* 批量添加文档
* 自动分批处理,避免大批量数据导致的内存溢出
*/
public BulkResult batchAdd(String indexName, List<?> dataList) {
ParamValidator.validateIndexName(indexName);
if (CollectionUtils.isEmpty(dataList)) {
return BulkResult.empty();
}
int batchSize = properties.getBulkBatchSize();
int totalSize = dataList.size();
int batchCount = (int) Math.ceil((double) totalSize / batchSize);
BulkResult.Builder resultBuilder = BulkResult.builder();
for (int i = 0; i < batchCount; i++) {
int fromIndex = i * batchSize;
int toIndex = Math.min(fromIndex + batchSize, totalSize);
List<?> batchData = dataList.subList(fromIndex, toIndex);
try {
BulkResult batchResult = doBatchAdd(indexName, batchData);
resultBuilder.merge(batchResult);
log.debug("批量添加进度: {}/{}, 成功: {}, 失败: {}",
toIndex, totalSize,
batchResult.getSuccessCount(),
batchResult.getFailureCount());
} catch (Exception e) {
log.error("批量添加第 {} 批失败", i + 1, e);
resultBuilder.addFailure(batchData.size(), e.getMessage());
}
}
BulkResult result = resultBuilder.build();
log.info("批量添加完成, index: {}, 总数: {}, 成功: {}, 失败: {}",
indexName, totalSize, result.getSuccessCount(), result.getFailureCount());
return result;
}
private BulkResult doBatchAdd(String indexName, List<?> dataList) throws IOException {
BulkRequest bulkRequest = new BulkRequest();
for (Object data : dataList) {
IndexRequest request = new IndexRequest(indexName);
request.source(JSON.toJSONString(data), XContentType.JSON);
if (properties.getEnableType()) {
request.type(properties.getDefaultType());
}
bulkRequest.add(request);
}
// 设置刷新策略
setRefreshPolicy(bulkRequest);
BulkResponse response = client.bulk(bulkRequest, RequestOptions.DEFAULT);
return parseBulkResponse(response);
}
/**
* 批量添加文档(带ID)
*/
public BulkResult batchAddWithId(String indexName, Map<String, Object> dataMap) {
ParamValidator.validateIndexName(indexName);
if (CollectionUtils.isEmpty(dataMap)) {
return BulkResult.empty();
}
int batchSize = properties.getBulkBatchSize();
List<Map.Entry<String, Object>> entries = new ArrayList<>(dataMap.entrySet());
int totalSize = entries.size();
int batchCount = (int) Math.ceil((double) totalSize / batchSize);
BulkResult.Builder resultBuilder = BulkResult.builder();
for (int i = 0; i < batchCount; i++) {
int fromIndex = i * batchSize;
int toIndex = Math.min(fromIndex + batchSize, totalSize);
List<Map.Entry<String, Object>> batchEntries = entries.subList(fromIndex, toIndex);
try {
BulkResult batchResult = doBatchAddWithId(indexName, batchEntries);
resultBuilder.merge(batchResult);
} catch (Exception e) {
log.error("批量添加第 {} 批失败", i + 1, e);
resultBuilder.addFailure(batchEntries.size(), e.getMessage());
}
}
BulkResult result = resultBuilder.build();
log.info("批量添加完成, index: {}, 总数: {}, 成功: {}, 失败: {}",
indexName, totalSize, result.getSuccessCount(), result.getFailureCount());
return result;
}
private BulkResult doBatchAddWithId(String indexName,
List<Map.Entry<String, Object>> entries)
throws IOException {
BulkRequest bulkRequest = new BulkRequest();
for (Map.Entry<String, Object> entry : entries) {
IndexRequest request = new IndexRequest(indexName);
request.id(entry.getKey());
request.source(JSON.toJSONString(entry.getValue()), XContentType.JSON);
if (properties.getEnableType()) {
request.type(properties.getDefaultType());
}
bulkRequest.add(request);
}
setRefreshPolicy(bulkRequest);
BulkResponse response = client.bulk(bulkRequest, RequestOptions.DEFAULT);
return parseBulkResponse(response);
}
// ==================== 查询操作 ====================
/**
* 根据ID查询
*/
public <T> Optional<T> getById(String indexName, String id, Class<T> clazz) {
ParamValidator.validateIndexName(indexName);
ParamValidator.validateDocumentId(id);
return RetryUtil.executeWithRetry(
() -> doGetById(indexName, id, clazz),
properties.getRetryTimes(),
properties.getRetryInterval(),
RetryUtil.defaultRetryCondition(),
"GET_BY_ID"
);
}
private <T> Optional<T> doGetById(String indexName, String id, Class<T> clazz)
throws IOException {
GetRequest request = new GetRequest(indexName, id);
GetResponse response = client.get(request, RequestOptions.DEFAULT);
if (!response.isExists()) {
return Optional.empty();
}
String source = response.getSourceAsString();
if (source == null || source.isEmpty()) {
return Optional.empty();
}
T obj = JSON.parseObject(source, clazz);
return Optional.of(obj);
}
/**
* 批量根据ID查询
*/
public <T> Map<String, T> getByIds(String indexName, List<String> ids, Class<T> clazz) {
ParamValidator.validateIndexName(indexName);
if (CollectionUtils.isEmpty(ids)) {
return Collections.emptyMap();
}
try {
MultiGetRequest request = new MultiGetRequest();
for (String id : ids) {
request.add(new MultiGetRequest.Item(indexName, id));
}
MultiGetResponse response = client.mget(request, RequestOptions.DEFAULT);
Map<String, T> result = new HashMap<>();
for (MultiGetItemResponse itemResponse : response.getResponses()) {
if (!itemResponse.isFailed() && itemResponse.getResponse().isExists()) {
String source = itemResponse.getResponse().getSourceAsString();
if (source != null && !source.isEmpty()) {
T obj = JSON.parseObject(source, clazz);
result.put(itemResponse.getId(), obj);
}
}
}
return result;
} catch (IOException e) {
throw new ElasticsearchException("GET_BY_IDS", indexName, e);
}
}
/**
* 判断文档是否存在
*/
public boolean exists(String indexName, String id) {
ParamValidator.validateIndexName(indexName);
ParamValidator.validateDocumentId(id);
try {
GetRequest request = new GetRequest(indexName, id);
request.fetchSourceContext(
org.elasticsearch.search.fetch.subphase.FetchSourceContext.DO_NOT_FETCH_SOURCE);
GetResponse response = client.get(request, RequestOptions.DEFAULT);
return response.isExists();
} catch (IOException e) {
throw new ElasticsearchException("EXISTS", indexName, e);
}
}
// ==================== 更新操作 ====================
/**
* 更新文档(全量更新,不存在则新增)
*/
public boolean upsert(String indexName, String id, Object data) {
ParamValidator.validateIndexName(indexName);
ParamValidator.validateDocumentId(id);
ParamValidator.notNull(data, "更新数据不能为空");
return RetryUtil.executeWithRetry(
() -> doUpsert(indexName, id, data),
properties.getRetryTimes(),
properties.getRetryInterval(),
RetryUtil.defaultRetryCondition(),
"UPSERT"
);
}
private boolean doUpsert(String indexName, String id, Object data) throws IOException {
UpdateRequest request = new UpdateRequest(indexName, id);
request.doc(JSON.toJSONString(data), XContentType.JSON);
request.docAsUpsert(true);
setRefreshPolicy(request);
UpdateResponse response = client.update(request, RequestOptions.DEFAULT);
log.info("更新文档成功, index: {}, id: {}, result: {}",
indexName, id, response.getResult());
return response.getResult() != DocWriteResponse.Result.NOOP;
}
/**
* 部分更新文档
*/
public boolean partialUpdate(String indexName, String id, Map<String, Object> fields) {
ParamValidator.validateIndexName(indexName);
ParamValidator.validateDocumentId(id);
if (CollectionUtils.isEmpty(fields)) {
return false;
}
return RetryUtil.executeWithRetry(
() -> doPartialUpdate(indexName, id, fields),
properties.getRetryTimes(),
properties.getRetryInterval(),
RetryUtil.defaultRetryCondition(),
"PARTIAL_UPDATE"
);
}
private boolean doPartialUpdate(String indexName, String id,
Map<String, Object> fields) throws IOException {
UpdateRequest request = new UpdateRequest(indexName, id);
request.doc(fields);
setRefreshPolicy(request);
UpdateResponse response = client.update(request, RequestOptions.DEFAULT);
log.info("部分更新文档成功, index: {}, id: {}", indexName, id);
return response.getResult() != DocWriteResponse.Result.NOOP;
}
/**
* 使用脚本更新
*/
public boolean updateByScript(String indexName, String id, String script,
Map<String, Object> params) {
ParamValidator.validateIndexName(indexName);
ParamValidator.validateDocumentId(id);
ParamValidator.notBlank(script, "脚本不能为空");
try {
UpdateRequest request = new UpdateRequest(indexName, id);
Script painlessScript = new Script(
ScriptType.INLINE,
"painless",
script,
params != null ? params : Collections.emptyMap()
);
request.script(painlessScript);
setRefreshPolicy(request);
UpdateResponse response = client.update(request, RequestOptions.DEFAULT);
log.info("脚本更新文档成功, index: {}, id: {}", indexName, id);
return response.getResult() != DocWriteResponse.Result.NOOP;
} catch (IOException e) {
throw new ElasticsearchException("UPDATE_BY_SCRIPT", indexName, e);
}
}
/**
* 根据条件批量更新
*/
public long updateByQuery(String indexName, QueryBuilder query,
String script, Map<String, Object> params) {
ParamValidator.validateIndexName(indexName);
ParamValidator.notNull(query, "查询条件不能为空");
ParamValidator.notBlank(script, "脚本不能为空");
try {
UpdateByQueryRequest request = new UpdateByQueryRequest(indexName);
request.setQuery(query);
request.setScript(new Script(
ScriptType.INLINE,
"painless",
script,
params != null ? params : Collections.emptyMap()
));
request.setRefresh(true);
BulkByScrollResponse response = client.updateByQuery(
request, RequestOptions.DEFAULT);
log.info("根据条件更新文档成功, index: {}, updated: {}",
indexName, response.getUpdated());
return response.getUpdated();
} catch (IOException e) {
throw new ElasticsearchException("UPDATE_BY_QUERY", indexName, e);
}
}
/**
* 批量更新
*/
public BulkResult batchUpdate(String indexName, Map<String, Object> dataMap) {
ParamValidator.validateIndexName(indexName);
if (CollectionUtils.isEmpty(dataMap)) {
return BulkResult.empty();
}
try {
BulkRequest bulkRequest = new BulkRequest();
dataMap.forEach((id, data) -> {
UpdateRequest request = new UpdateRequest(indexName, id);
request.doc(JSON.toJSONString(data), XContentType.JSON);
request.docAsUpsert(true);
bulkRequest.add(request);
});
setRefreshPolicy(bulkRequest);
BulkResponse response = client.bulk(bulkRequest, RequestOptions.DEFAULT);
BulkResult result = parseBulkResponse(response);
log.info("批量更新完成, index: {}, 成功: {}, 失败: {}",
indexName, result.getSuccessCount(), result.getFailureCount());
return result;
} catch (IOException e) {
throw new ElasticsearchException("BATCH_UPDATE", indexName, e);
}
}
// ==================== 删除操作 ====================
/**
* 删除文档
*/
public boolean delete(String indexName, String id) {
ParamValidator.validateIndexName(indexName);
ParamValidator.validateDocumentId(id);
return RetryUtil.executeWithRetry(
() -> doDelete(indexName, id),
properties.getRetryTimes(),
properties.getRetryInterval(),
RetryUtil.defaultRetryCondition(),
"DELETE"
);
}
private boolean doDelete(String indexName, String id) throws IOException {
DeleteRequest request = new DeleteRequest(indexName, id);
setRefreshPolicy(request);
DeleteResponse response = client.delete(request, RequestOptions.DEFAULT);
log.info("删除文档成功, index: {}, id: {}", indexName, id);
return response.getResult() == DocWriteResponse.Result.DELETED;
}
/**
* 批量删除
*/
public BulkResult batchDelete(String indexName, List<String> ids) {
ParamValidator.validateIndexName(indexName);
if (CollectionUtils.isEmpty(ids)) {
return BulkResult.empty();
}
try {
BulkRequest bulkRequest = new BulkRequest();
for (String id : ids) {
DeleteRequest request = new DeleteRequest(indexName, id);
bulkRequest.add(request);
}
setRefreshPolicy(bulkRequest);
BulkResponse response = client.bulk(bulkRequest, RequestOptions.DEFAULT);
BulkResult result = parseBulkResponse(response);
log.info("批量删除完成, index: {}, 成功: {}, 失败: {}",
indexName, result.getSuccessCount(), result.getFailureCount());
return result;
} catch (IOException e) {
throw new ElasticsearchException("BATCH_DELETE", indexName, e);
}
}
/**
* 根据条件删除
*/
public long deleteByQuery(String indexName, QueryBuilder query) {
ParamValidator.validateIndexName(indexName);
ParamValidator.notNull(query, "查询条件不能为空");
try {
DeleteByQueryRequest request = new DeleteByQueryRequest(indexName);
request.setQuery(query);
request.setRefresh(true);
BulkByScrollResponse response = client.deleteByQuery(
request, RequestOptions.DEFAULT);
log.info("根据条件删除文档成功, index: {}, deleted: {}",
indexName, response.getDeleted());
return response.getDeleted();
} catch (IOException e) {
throw new ElasticsearchException("DELETE_BY_QUERY", indexName, e);
}
}
// ==================== 内部工具方法 ====================
/**
* 设置刷新策略
*/
private void setRefreshPolicy(IndexRequest request) {
String policy = properties.getBulkRefreshPolicy();
if ("immediate".equals(policy)) {
request.setRefreshPolicy(WriteRequest.RefreshPolicy.IMMEDIATE);
} else if ("wait_for".equals(policy)) {
request.setRefreshPolicy(WriteRequest.RefreshPolicy.WAIT_UNTIL);
}
}
private void setRefreshPolicy(UpdateRequest request) {
String policy = properties.getBulkRefreshPolicy();
if ("immediate".equals(policy)) {
request.setRefreshPolicy(WriteRequest.RefreshPolicy.IMMEDIATE);
} else if ("wait_for".equals(policy)) {
request.setRefreshPolicy(WriteRequest.RefreshPolicy.WAIT_UNTIL);
}
}
private void setRefreshPolicy(DeleteRequest request) {
String policy = properties.getBulkRefreshPolicy();
if ("immediate".equals(policy)) {
request.setRefreshPolicy(WriteRequest.RefreshPolicy.IMMEDIATE);
} else if ("wait_for".equals(policy)) {
request.setRefreshPolicy(WriteRequest.RefreshPolicy.WAIT_UNTIL);
}
}
private void setRefreshPolicy(BulkRequest request) {
String policy = properties.getBulkRefreshPolicy();
if ("immediate".equals(policy)) {
request.setRefreshPolicy(WriteRequest.RefreshPolicy.IMMEDIATE);
} else if ("wait_for".equals(policy)) {
request.setRefreshPolicy(WriteRequest.RefreshPolicy.WAIT_UNTIL);
}
}
/**
* 解析批量操作响应
*/
private BulkResult parseBulkResponse(BulkResponse response) {
BulkResult.Builder builder = BulkResult.builder();
for (BulkItemResponse itemResponse : response.getItems()) {
if (itemResponse.isFailed()) {
builder.addFailure(itemResponse.getId(),
itemResponse.getFailureMessage());
} else {
builder.addSuccess(itemResponse.getId());
}
}
return builder.build();
}
// ==================== 批量操作结果 ====================
/**
* 批量操作结果
*/
@lombok.Data
public static class BulkResult {
private int successCount;
private int failureCount;
private List<String> successIds;
private List<FailureItem> failures;
@lombok.Data
@lombok.AllArgsConstructor
public static class FailureItem {
private String id;
private String reason;
}
public static BulkResult empty() {
BulkResult result = new BulkResult();
result.successCount = 0;
result.failureCount = 0;
result.successIds = Collections.emptyList();
result.failures = Collections.emptyList();
return result;
}
public boolean hasFailures() {
return failureCount > 0;
}
public boolean isAllSuccess() {
return failureCount == 0;
}
public static Builder builder() {
return new Builder();
}
public static class Builder {
private List<String> successIds = new ArrayList<>();
private List<FailureItem> failures = new ArrayList<>();
public Builder addSuccess(String id) {
successIds.add(id);
return this;
}
public Builder addFailure(String id, String reason) {
failures.add(new FailureItem(id, reason));
return this;
}
public Builder addFailure(int count, String reason) {
for (int i = 0; i < count; i++) {
failures.add(new FailureItem(null, reason));
}
return this;
}
public Builder merge(BulkResult other) {
if (other.successIds != null) {
successIds.addAll(other.successIds);
}
if (other.failures != null) {
failures.addAll(other.failures);
}
return this;
}
public BulkResult build() {
BulkResult result = new BulkResult();
result.successIds = successIds;
result.failures = failures;
result.successCount = successIds.size();
result.failureCount = failures.size();
return result;
}
}
}
}
9.12 滚动查询结果类
package com.example.elasticsearch.result;
import lombok.AllArgsConstructor;
import lombok.Data;
import lombok.NoArgsConstructor;
import java.util.List;
/**
* 滚动查询结果
*/
@Data
@NoArgsConstructor
@AllArgsConstructor
public class ScrollResult<T> {
/**
* 滚动ID(用于获取下一批数据)
*/
private String scrollId;
/**
* 当前批次数据
*/
private List<T> records;
/**
* 总记录数
*/
private Long total;
/**
* 是否是最后一批(没有更多数据)
*/
private Boolean finished;
/**
* 判断是否还有更多数据
*/
public boolean hasMore() {
return !finished && records != null && !records.isEmpty();
}
}
9.13 索引操作服务
package com.example.elasticsearch.core;
import com.example.elasticsearch.config.ElasticsearchProperties;
import com.example.elasticsearch.exception.ElasticsearchException;
import lombok.extern.slf4j.Slf4j;
import org.elasticsearch.action.admin.indices.delete.DeleteIndexRequest;
import org.elasticsearch.action.support.master.AcknowledgedResponse;
import org.elasticsearch.client.RequestOptions;
import org.elasticsearch.client.RestHighLevelClient;
import org.elasticsearch.client.indices.*;
import org.elasticsearch.common.settings.Settings;
import org.elasticsearch.common.xcontent.XContentType;
import java.io.IOException;
import java.util.Map;
/**
* Elasticsearch 索引操作服务
*/
@Slf4j
public class ElasticsearchIndexService {
private final RestHighLevelClient client;
private final ElasticsearchProperties properties;
public ElasticsearchIndexService(RestHighLevelClient client,
ElasticsearchProperties properties) {
this.client = client;
this.properties = properties;
}
/**
* 创建索引
*/
public boolean createIndex(String indexName) {
return createIndex(indexName, null, null);
}
/**
* 创建索引(带 Mapping)
*/
public boolean createIndex(String indexName, String mapping) {
try {
CreateIndexRequest request = new CreateIndexRequest(indexName);
if (mapping != null) {
request.source(mapping, XContentType.JSON);
}
CreateIndexResponse response = client.indices()
.create(request, RequestOptions.DEFAULT);
log.info("创建索引成功, index: {}", indexName);
return response.isAcknowledged();
} catch (IOException e) {
log.error("创建索引失败", e);
throw new ElasticsearchException("CREATE_INDEX", indexName, e);
}
}
/**
* 创建索引(带 Settings 和 Mapping)
*/
public boolean createIndex(String indexName,
Map<String, Object> settings,
Map<String, Object> mapping) {
try {
CreateIndexRequest request = new CreateIndexRequest(indexName);
if (settings != null) {
request.settings(settings);
}
if (mapping != null) {
request.mapping(mapping);
}
CreateIndexResponse response = client.indices()
.create(request, RequestOptions.DEFAULT);
log.info("创建索引成功, index: {}", indexName);
return response.isAcknowledged();
} catch (IOException e) {
log.error("创建索引失败", e);
throw new ElasticsearchException("CREATE_INDEX", indexName, e);
}
}
/**
* 判断索引是否存在
*/
public boolean existsIndex(String indexName) {
try {
GetIndexRequest request = new GetIndexRequest(indexName);
return client.indices().exists(request, RequestOptions.DEFAULT);
} catch (IOException e) {
log.error("判断索引是否存在失败", e);
throw new ElasticsearchException("EXISTS_INDEX", indexName, e);
}
}
/**
* 删除索引
*/
public boolean deleteIndex(String indexName) {
try {
DeleteIndexRequest request = new DeleteIndexRequest(indexName);
AcknowledgedResponse response = client.indices()
.delete(request, RequestOptions.DEFAULT);
log.info("删除索引成功, index: {}", indexName);
return response.isAcknowledged();
} catch (IOException e) {
log.error("删除索引失败", e);
throw new ElasticsearchException("DELETE_INDEX", indexName, e);
}
}
/**
* 更新 Mapping
*/
public boolean updateMapping(String indexName, Map<String, Object> properties) {
try {
PutMappingRequest request = new PutMappingRequest(indexName);
request.source(Map.of("properties", properties));
AcknowledgedResponse response = client.indices()
.putMapping(request, RequestOptions.DEFAULT);
log.info("更新 Mapping 成功, index: {}", indexName);
return response.isAcknowledged();
} catch (IOException e) {
log.error("更新 Mapping 失败", e);
throw new ElasticsearchException("UPDATE_MAPPING", indexName, e);
}
}
/**
* 获取索引信息
*/
public GetIndexResponse getIndex(String indexName) {
try {
GetIndexRequest request = new GetIndexRequest(indexName);
return client.indices().get(request, RequestOptions.DEFAULT);
} catch (IOException e) {
log.error("获取索引信息失败", e);
throw new ElasticsearchException("GET_INDEX", indexName, e);
}
}
/**
* 刷新索引
*/
public void refreshIndex(String... indexNames) {
try {
org.elasticsearch.client.indices.RefreshRequest request =
new org.elasticsearch.client.indices.RefreshRequest(indexNames);
client.indices().refresh(request, RequestOptions.DEFAULT);
log.info("刷新索引成功, indices: {}", String.join(",", indexNames));
} catch (IOException e) {
log.error("刷新索引失败", e);
throw new ElasticsearchException("REFRESH_INDEX",
String.join(",", indexNames), e);
}
}
/**
* 创建索引别名
*/
public boolean createAlias(String indexName, String aliasName) {
try {
var request = new org.elasticsearch.client.indices.IndicesAliasesRequest();
var action = new org.elasticsearch.client.indices.IndicesAliasesRequest
.AliasActions(org.elasticsearch.client.indices.IndicesAliasesRequest
.AliasActions.Type.ADD)
.index(indexName)
.alias(aliasName);
request.addAliasAction(action);
var response = client.indices().updateAliases(request, RequestOptions.DEFAULT);
log.info("创建索引别名成功, index: {}, alias: {}", indexName, aliasName);
return response.isAcknowledged();
} catch (IOException e) {
log.error("创建索引别名失败", e);
throw new ElasticsearchException("CREATE_ALIAS", indexName, e);
}
}
}
9.14 架构图
- ✅ 链式调用 -
EsQueryBuilder支持流畅的链式构建 - ✅ 类型安全 - 泛型支持,自动转换结果类型
- ✅ 功能完整 - 覆盖索引、文档、查询、聚合操作
- ✅ 高亮支持 - 内置高亮配置和结果解析
- ✅ 分页封装 - 统一的分页结果对象
- ✅ 连接池 - 支持集群、认证、连接池配置
- ✅ 版本兼容 - 支持 ES 6.x/7.x type 兼容
- ✅ 异常处理 - 统一的异常封装
十、完整使用文档
快速开始
添加依赖
<!-- pom.xml -->
<dependencies>
<!-- ES Starter -->
<dependency>
<groupId>com.example</groupId>
<artifactId>elasticsearch-spring-boot-starter</artifactId>
<version>1.0.0</version>
</dependency>
<!-- 或者使用官方依赖 -->
<dependency>
<groupId>org.springframework.boot</groupId>
<artifactId>spring-boot-starter-data-elasticsearch</artifactId>
</dependency>
<!-- FastJSON -->
<dependency>
<groupId>com.alibaba</groupId>
<artifactId>fastjson</artifactId>
<version>1.2.83</version>
</dependency>
</dependencies>
配置文件
# application.yml
elasticsearch:
enabled: true
nodes:
- localhost:9200
username: elastic # 可选
password: your_password # 可选
scheme: http
# 连接配置
connect-timeout: 5000
socket-timeout: 30000
max-conn-total: 100
max-conn-per-route: 50
# 重试配置
retry-times: 3
retry-interval: 1000
# 批量操作配置
bulk-batch-size: 1000
bulk-refresh-policy: none # none, immediate, wait_for
# 日志配置
print-dsl: true # 开发环境建议开启
slow-query-threshold: 3000 # 慢查询阈值(ms)
# 版本兼容
enable-type: false # ES 7.x 设为 false
创建实体类
package com.example.demo.entity;
import lombok.Data;
import java.util.Date;
/**
* 节目实体
*/
@Data
public class Program {
private Long id;
/**
* 节目标题
*/
private String title;
/**
* 演员
*/
private String actor;
/**
* 分类ID
*/
private Long categoryId;
/**
* 分类名称
*/
private String categoryName;
/**
* 地区ID
*/
private Long areaId;
/**
* 地区名称
*/
private String areaName;
/**
* 演出时间
*/
private Date showTime;
/**
* 最低价格
*/
private Double minPrice;
/**
* 最高价格
*/
private Double maxPrice;
/**
* 状态:1-上架, 0-下架
*/
private Integer status;
/**
* 封面图片
*/
private String coverImage;
/**
* 创建时间
*/
private Date createTime;
/**
* 更新时间
*/
private Date updateTime;
}
创建搜索参数类
package com.example.demo.dto;
import lombok.Data;
import java.util.Date;
/**
* 节目搜索参数
*/
@Data
public class ProgramSearchParam {
/**
* 关键词(搜索标题和演员)
*/
private String keyword;
/**
* 分类ID
*/
private Long categoryId;
/**
* 地区ID
*/
private Long areaId;
/**
* 状态
*/
private Integer status;
/**
* 最低价格
*/
private Double minPrice;
/**
* 最高价格
*/
private Double maxPrice;
/**
* 开始时间
*/
private Date startTime;
/**
* 结束时间
*/
private Date endTime;
/**
* 页码
*/
private Integer pageNum = 1;
/**
* 每页大小
*/
private Integer pageSize = 10;
/**
* 排序字段
*/
private String sortField = "showTime";
/**
* 排序方式:asc, desc
*/
private String sortOrder = "asc";
}
完整业务示例
package com.example.demo.service;
import com.example.demo.dto.ProgramSearchParam;
import com.example.demo.entity.Program;
import com.example.elasticsearch.core.ElasticsearchDocumentService;
import com.example.elasticsearch.core.ElasticsearchDocumentService.BulkResult;
import com.example.elasticsearch.core.ElasticsearchIndexService;
import com.example.elasticsearch.core.ElasticsearchService;
import com.example.elasticsearch.query.EsQueryBuilder;
import com.example.elasticsearch.result.AggregationResult;
import com.example.elasticsearch.result.PageResult;
import com.example.elasticsearch.result.ScrollResult;
import com.example.elasticsearch.result.SearchResult;
import lombok.RequiredArgsConstructor;
import lombok.extern.slf4j.Slf4j;
import org.elasticsearch.search.sort.SortOrder;
import org.springframework.stereotype.Service;
import javax.annotation.PostConstruct;
import java.util.*;
import java.util.function.Consumer;
import java.util.stream.Collectors;
/**
* 节目搜索服务 - 完整示例
*/
@Slf4j
@Service
@RequiredArgsConstructor
public class ProgramSearchService {
private final ElasticsearchService esService;
private final ElasticsearchDocumentService documentService;
private final ElasticsearchIndexService indexService;
private static final String INDEX_NAME = "program";
// ==================== 索引管理 ====================
/**
* 初始化索引(应用启动时调用)
*/
@PostConstruct
public void initIndex() {
if (!indexService.existsIndex(INDEX_NAME)) {
String mapping = """
{
"mappings": {
"properties": {
"id": { "type": "long" },
"title": {
"type": "text",
"analyzer": "ik_max_word",
"search_analyzer": "ik_smart",
"fields": {
"keyword": { "type": "keyword" }
}
},
"actor": {
"type": "text",
"analyzer": "ik_max_word",
"fields": {
"keyword": { "type": "keyword" }
}
},
"categoryId": { "type": "long" },
"categoryName": { "type": "keyword" },
"areaId": { "type": "long" },
"areaName": { "type": "keyword" },
"showTime": { "type": "date" },
"minPrice": { "type": "double" },
"maxPrice": { "type": "double" },
"status": { "type": "integer" },
"coverImage": { "type": "keyword", "index": false },
"createTime": { "type": "date" },
"updateTime": { "type": "date" }
}
},
"settings": {
"number_of_shards": 3,
"number_of_replicas": 1,
"refresh_interval": "1s"
}
}
""";
indexService.createIndex(INDEX_NAME, mapping);
log.info("索引 [{}] 创建成功", INDEX_NAME);
}
}
/**
* 重建索引
*/
public void rebuildIndex() {
// 删除旧索引
if (indexService.existsIndex(INDEX_NAME)) {
indexService.deleteIndex(INDEX_NAME);
}
// 重新初始化
initIndex();
}
// ==================== 文档操作 ====================
/**
* 添加/更新节目
*/
public void saveProgram(Program program) {
program.setUpdateTime(new Date());
if (program.getCreateTime() == null) {
program.setCreateTime(new Date());
}
documentService.upsert(INDEX_NAME, String.valueOf(program.getId()), program);
}
/**
* 批量添加节目
*/
public BulkResult batchSavePrograms(List<Program> programs) {
Date now = new Date();
programs.forEach(p -> {
p.setUpdateTime(now);
if (p.getCreateTime() == null) {
p.setCreateTime(now);
}
});
return documentService.batchAdd(INDEX_NAME, programs);
}
/**
* 批量添加节目(带ID)
*/
public BulkResult batchSaveProgramsWithId(List<Program> programs) {
Date now = new Date();
Map<String, Object> dataMap = new HashMap<>();
programs.forEach(p -> {
p.setUpdateTime(now);
if (p.getCreateTime() == null) {
p.setCreateTime(now);
}
dataMap.put(String.valueOf(p.getId()), p);
});
return documentService.batchAddWithId(INDEX_NAME, dataMap);
}
/**
* 根据ID查询节目
*/
public Optional<Program> getProgramById(Long id) {
return documentService.getById(INDEX_NAME, String.valueOf(id), Program.class);
}
/**
* 批量查询节目
*/
public Map<Long, Program> getProgramsByIds(List<Long> ids) {
List<String> strIds = ids.stream()
.map(String::valueOf)
.collect(Collectors.toList());
Map<String, Program> resultMap = documentService.getByIds(
INDEX_NAME, strIds, Program.class);
// 转换 key 类型
return resultMap.entrySet().stream()
.collect(Collectors.toMap(
e -> Long.parseLong(e.getKey()),
Map.Entry::getValue
));
}
/**
* 删除节目
*/
public void deleteProgram(Long id) {
documentService.delete(INDEX_NAME, String.valueOf(id));
}
/**
* 批量删除节目
*/
public BulkResult batchDeletePrograms(List<Long> ids) {
List<String> strIds = ids.stream()
.map(String::valueOf)
.collect(Collectors.toList());
return documentService.batchDelete(INDEX_NAME, strIds);
}
/**
* 更新节目状态
*/
public boolean updateProgramStatus(Long id, Integer status) {
Map<String, Object> fields = new HashMap<>();
fields.put("status", status);
fields.put("updateTime", new Date());
return documentService.partialUpdate(INDEX_NAME, String.valueOf(id), fields);
}
/**
* 批量更新节目状态(使用脚本)
*/
public long batchUpdateStatus(List<Long> ids, Integer status) {
EsQueryBuilder queryBuilder = EsQueryBuilder.builder(INDEX_NAME)
.filterTerms("id", ids);
String script = "ctx._source.status = params.status; " +
"ctx._source.updateTime = params.updateTime";
Map<String, Object> params = new HashMap<>();
params.put("status", status);
params.put("updateTime", System.currentTimeMillis());
return documentService.updateByQuery(
INDEX_NAME, queryBuilder.getBoolQuery(), script, params);
}
// ==================== 基础搜索 ====================
/**
* 根据分类查询节目列表
*/
public List<Program> findByCategory(Long categoryId) {
return esService.query(INDEX_NAME, "categoryId", categoryId, Program.class);
}
/**
* 根据多条件查询
*/
public List<Program> findByConditions(Long categoryId, Long areaId, Integer status) {
Map<String, Object> params = new HashMap<>();
params.put("categoryId", categoryId);
params.put("areaId", areaId);
params.put("status", status);
return esService.query(INDEX_NAME, params, Program.class);
}
/**
* 简单分页查询
*/
public PageResult<Program> findPage(int pageNum, int pageSize) {
return esService.queryPage(INDEX_NAME, null, pageNum, pageSize, Program.class);
}
// ==================== 高级搜索 ====================
/**
* 关键词搜索(标题 + 演员)
*/
public List<Program> searchByKeyword(String keyword) {
EsQueryBuilder queryBuilder = EsQueryBuilder.builder(INDEX_NAME)
.shouldMatch(keyword, "title", "actor")
.filterTerm("status", 1) // 只搜索上架的
.sort("showTime", SortOrder.ASC);
return esService.search(queryBuilder, Program.class);
}
/**
* 关键词搜索(带高亮)
*/
public List<SearchResult<Program>> searchWithHighlight(String keyword) {
EsQueryBuilder queryBuilder = EsQueryBuilder.builder(INDEX_NAME)
.shouldMatch(keyword, "title", "actor")
.filterTerm("status", 1)
.highlight("<span class='highlight'>", "</span>", "title", "actor")
.sort("showTime", SortOrder.ASC);
return esService.searchWithHighlight(queryBuilder, Program.class);
}
/**
* 复杂条件搜索 + 分页
*/
public PageResult<Program> search(ProgramSearchParam param) {
EsQueryBuilder queryBuilder = EsQueryBuilder.builder(INDEX_NAME)
// 关键词搜索(标题或演员)
.shouldMatch(param.getKeyword(), "title", "actor")
// 分类过滤
.filterTerm("categoryId", param.getCategoryId())
// 地区过滤
.filterTerm("areaId", param.getAreaId())
// 状态过滤
.filterTerm("status", param.getStatus())
// 价格范围
.filterRange("minPrice", param.getMinPrice(), param.getMaxPrice())
// 时间范围
.filterRange("showTime",
param.getStartTime() != null ? param.getStartTime().getTime() : null,
param.getEndTime() != null ? param.getEndTime().getTime() : null)
// 排序
.sort(param.getSortField(),
"desc".equalsIgnoreCase(param.getSortOrder())
? SortOrder.DESC : SortOrder.ASC);
return esService.searchPage(queryBuilder,
param.getPageNum(), param.getPageSize(), Program.class);
}
/**
* 复杂条件搜索 + 分页 + 高亮
*/
public PageResult<SearchResult<Program>> searchWithHighlight(ProgramSearchParam param) {
EsQueryBuilder queryBuilder = EsQueryBuilder.builder(INDEX_NAME)
.shouldMatch(param.getKeyword(), "title", "actor")
.filterTerm("categoryId", param.getCategoryId())
.filterTerm("areaId", param.getAreaId())
.filterTerm("status", param.getStatus())
.filterRange("minPrice", param.getMinPrice(), param.getMaxPrice())
.filterRange("showTime",
param.getStartTime() != null ? param.getStartTime().getTime() : null,
param.getEndTime() != null ? param.getEndTime().getTime() : null)
.highlight("title", "actor")
.sort(param.getSortField(),
"desc".equalsIgnoreCase(param.getSortOrder())
? SortOrder.DESC : SortOrder.ASC);
return esService.searchPageWithHighlight(queryBuilder,
param.getPageNum(), param.getPageSize(), Program.class);
}
// ==================== 深度分页(Search After) ====================
/**
* 使用 Search After 进行深度分页
* 适用于页码超过 1000 的场景
*
* @param param 搜索参数
* @param searchAfterValue 上一页最后一条记录的排序值
* @return 搜索结果(包含 sortValues 用于下一页查询)
*/
public List<SearchResult<Program>> searchAfter(ProgramSearchParam param,
Object[] searchAfterValue) {
EsQueryBuilder queryBuilder = EsQueryBuilder.builder(INDEX_NAME)
.shouldMatch(param.getKeyword(), "title", "actor")
.filterTerm("categoryId", param.getCategoryId())
.filterTerm("status", param.getStatus())
// 必须有排序字段,且最后一个字段应该是唯一的(如 _id)
.sort(param.getSortField(),
"desc".equalsIgnoreCase(param.getSortOrder())
? SortOrder.DESC : SortOrder.ASC)
.sort("id", SortOrder.ASC); // 使用 id 保证唯一性
return esService.searchAfter(queryBuilder, searchAfterValue,
param.getPageSize(), Program.class);
}
// ==================== 滚动查询(大数据导出) ====================
/**
* 滚动查询导出所有数据
* 适用于数据导出场景
*
* @param param 搜索参数
* @param consumer 数据消费者
*/
public void scrollExport(ProgramSearchParam param, Consumer<List<Program>> consumer) {
EsQueryBuilder queryBuilder = EsQueryBuilder.builder(INDEX_NAME)
.shouldMatch(param.getKeyword(), "title", "actor")
.filterTerm("categoryId", param.getCategoryId())
.filterTerm("status", param.getStatus())
.sort("id", SortOrder.ASC); // 滚动查询必须有排序
int batchSize = 1000;
int scrollTime = 5; // 滚动上下文保持时间(分钟)
// 初始化滚动查询
ScrollResult<Program> scrollResult = esService.scrollSearch(
queryBuilder, scrollTime, batchSize, Program.class);
try {
// 处理首批数据
if (!scrollResult.getRecords().isEmpty()) {
consumer.accept(scrollResult.getRecords());
}
// 继续获取后续数据
while (scrollResult.hasMore()) {
scrollResult = esService.scrollNext(
scrollResult.getScrollId(), scrollTime, batchSize, Program.class);
if (!scrollResult.getRecords().isEmpty()) {
consumer.accept(scrollResult.getRecords());
}
}
log.info("滚动导出完成, 总数: {}", scrollResult.getTotal());
} finally {
// 清除滚动上下文
esService.clearScroll(scrollResult.getScrollId());
}
}
// ==================== 聚合统计 ====================
/**
* 按分类统计节目数量
*/
public Map<String, Long> countByCategory() {
EsQueryBuilder queryBuilder = EsQueryBuilder.builder(INDEX_NAME)
.filterTerm("status", 1) // 只统计上架的
.termsAggregation("category_count", "categoryId", 100);
Map<String, AggregationResult> aggResult = esService.searchAggregation(queryBuilder);
AggregationResult categoryAgg = aggResult.get("category_count");
if (categoryAgg == null || categoryAgg.getBuckets() == null) {
return Collections.emptyMap();
}
return categoryAgg.getBuckets().stream()
.collect(Collectors.toMap(
AggregationResult.BucketData::getKey,
AggregationResult.BucketData::getDocCount
));
}
/**
* 统计价格区间
*/
public Map<String, Double> getPriceStats() {
EsQueryBuilder queryBuilder = EsQueryBuilder.builder(INDEX_NAME)
.filterTerm("status", 1)
.minAggregation("min_price", "minPrice")
.maxAggregation("max_price", "maxPrice")
.avgAggregation("avg_price", "minPrice");
Map<String, AggregationResult> aggResult = esService.searchAggregation(queryBuilder);
Map<String, Double> stats = new HashMap<>();
if (aggResult.containsKey("min_price")) {
stats.put("minPrice", aggResult.get("min_price").getValue());
}
if (aggResult.containsKey("max_price")) {
stats.put("maxPrice", aggResult.get("max_price").getValue());
}
if (aggResult.containsKey("avg_price")) {
stats.put("avgPrice", aggResult.get("avg_price").getValue());
}
return stats;
}
/**
* 综合统计
*/
public Map<String, Object> getStatistics() {
EsQueryBuilder queryBuilder = EsQueryBuilder.builder(INDEX_NAME)
.filterTerm("status", 1)
.termsAggregation("by_category", "categoryId", 50)
.termsAggregation("by_area", "areaId", 50)
.avgAggregation("avg_price", "minPrice")
.minAggregation("min_price", "minPrice")
.maxAggregation("max_price", "maxPrice");
Map<String, AggregationResult> aggResult = esService.searchAggregation(queryBuilder);
Map<String, Object> stats = new HashMap<>();
stats.put("aggregations", aggResult);
stats.put("total", esService.count(
EsQueryBuilder.builder(INDEX_NAME).filterTerm("status", 1)));
return stats;
}
// ==================== 工具方法 ====================
/**
* 统计符合条件的数量
*/
public long count(ProgramSearchParam param) {
EsQueryBuilder queryBuilder = EsQueryBuilder.builder(INDEX_NAME)
.shouldMatch(param.getKeyword(), "title", "actor")
.filterTerm("categoryId", param.getCategoryId())
.filterTerm("areaId", param.getAreaId())
.filterTerm("status", param.getStatus());
return esService.count(queryBuilder);
}
/**
* 判断是否存在符合条件的数据
*/
public boolean exists(ProgramSearchParam param) {
return count(param) > 0;
}
}
Controller 示例
package com.example.demo.controller;
import com.example.demo.dto.ProgramSearchParam;
import com.example.demo.entity.Program;
import com.example.demo.service.ProgramSearchService;
import com.example.elasticsearch.core.ElasticsearchDocumentService.BulkResult;
import com.example.elasticsearch.result.PageResult;
import com.example.elasticsearch.result.SearchResult;
import lombok.RequiredArgsConstructor;
import org.springframework.web.bind.annotation.*;
import java.util.List;
import java.util.Map;
import java.util.Optional;
/**
* 节目搜索 API
*/
@RestController
@RequestMapping("/api/programs")
@RequiredArgsConstructor
public class ProgramController {
private final ProgramSearchService searchService;
// ==================== 基础 CRUD ====================
/**
* 添加/更新节目
*/
@PostMapping
public String saveProgram(@RequestBody Program program) {
searchService.saveProgram(program);
return "success";
}
/**
* 批量添加节目
*/
@PostMapping("/batch")
public BulkResult batchSavePrograms(@RequestBody List<Program> programs) {
return searchService.batchSavePrograms(programs);
}
/**
* 根据ID查询
*/
@GetMapping("/{id}")
public Optional<Program> getProgram(@PathVariable Long id) {
return searchService.getProgramById(id);
}
/**
* 批量查询
*/
@PostMapping("/batch-get")
public Map<Long, Program> batchGetPrograms(@RequestBody List<Long> ids) {
return searchService.getProgramsByIds(ids);
}
/**
* 删除节目
*/
@DeleteMapping("/{id}")
public String deleteProgram(@PathVariable Long id) {
searchService.deleteProgram(id);
return "success";
}
/**
* 更新状态
*/
@PutMapping("/{id}/status")
public boolean updateStatus(@PathVariable Long id, @RequestParam Integer status) {
return searchService.updateProgramStatus(id, status);
}
// ==================== 搜索 ====================
/**
* 关键词搜索
*/
@GetMapping("/search")
public List<Program> search(@RequestParam String keyword) {
return searchService.searchByKeyword(keyword);
}
/**
* 关键词搜索(带高亮)
*/
@GetMapping("/search/highlight")
public List<SearchResult<Program>> searchWithHighlight(@RequestParam String keyword) {
return searchService.searchWithHighlight(keyword);
}
/**
* 高级搜索(分页)
*/
@PostMapping("/search/page")
public PageResult<Program> searchPage(@RequestBody ProgramSearchParam param) {
return searchService.search(param);
}
/**
* 高级搜索(分页 + 高亮)
*/
@PostMapping("/search/page-highlight")
public PageResult<SearchResult<Program>> searchPageWithHighlight(
@RequestBody ProgramSearchParam param) {
return searchService.searchWithHighlight(param);
}
/**
* 深度分页(Search After)
*/
@PostMapping("/search/after")
public List<SearchResult<Program>> searchAfter(
@RequestBody ProgramSearchParam param,
@RequestParam(required = false) Long afterId,
@RequestParam(required = false) Long afterShowTime) {
Object[] searchAfter = null;
if (afterShowTime != null && afterId != null) {
searchAfter = new Object[]{afterShowTime, afterId};
}
return searchService.searchAfter(param, searchAfter);
}
// ==================== 统计 ====================
/**
* 按分类统计
*/
@GetMapping("/stats/category")
public Map<String, Long> countByCategory() {
return searchService.countByCategory();
}
/**
* 价格统计
*/
@GetMapping("/stats/price")
public Map<String, Double> getPriceStats() {
return searchService.getPriceStats();
}
/**
* 综合统计
*/
@GetMapping("/stats")
public Map<String, Object> getStatistics() {
return searchService.getStatistics();
}
// ==================== 索引管理 ====================
/**
* 重建索引
*/
@PostMapping("/index/rebuild")
public String rebuildIndex() {
searchService.rebuildIndex();
return "success";
}
}
最佳实践指南
索引设计
/**
* 索引设计最佳实践
*/
public class IndexDesignBestPractice {
/**
* 1. 合理设置分片数
* - 单个分片大小建议 10-50GB
* - 分片数 = 数据量 / 单分片大小
*/
public static final int SHARD_COUNT = 3;
/**
* 2. 副本数设置
* - 生产环境至少 1 个副本
* - 写入密集型可先设为 0,完成后再设为 1
*/
public static final int REPLICA_COUNT = 1;
/**
* 3. Mapping 设计原则
* - 明确字段类型,避免动态映射
* - text 用于全文搜索,keyword 用于精确匹配/聚合
* - 不需要搜索的字段设置 index: false
* - 大文本字段考虑使用 store: true
*/
public static String createOptimizedMapping() {
return """
{
"mappings": {
"properties": {
"id": { "type": "long" },
"title": {
"type": "text",
"analyzer": "ik_max_word",
"search_analyzer": "ik_smart",
"fields": {
"keyword": {
"type": "keyword",
"ignore_above": 256
}
}
},
"description": {
"type": "text",
"analyzer": "ik_max_word",
"index": true,
"store": true
},
"status": { "type": "keyword" },
"price": { "type": "scaled_float", "scaling_factor": 100 },
"createTime": { "type": "date", "format": "epoch_millis" },
"tags": { "type": "keyword" },
"coverImage": { "type": "keyword", "index": false }
}
},
"settings": {
"number_of_shards": 3,
"number_of_replicas": 1,
"refresh_interval": "1s",
"max_result_window": 10000,
"analysis": {
"analyzer": {
"ik_smart_pinyin": {
"type": "custom",
"tokenizer": "ik_smart",
"filter": ["lowercase", "pinyin_filter"]
}
}
}
}
}
""";
}
}
查询优化
/**
* 查询优化最佳实践
*/
@Slf4j
@Service
public class QueryOptimizationService {
@Autowired
private ElasticsearchService esService;
private static final String INDEX_NAME = "program";
/**
* 1. 使用 filter 而非 must 进行过滤
* filter 不计算评分,性能更好,且结果可缓存
*/
public List<Program> optimizedFilterQuery(Long categoryId, Integer status) {
EsQueryBuilder queryBuilder = EsQueryBuilder.builder(INDEX_NAME)
// ✅ 正确:使用 filter 进行精确匹配
.filterTerm("categoryId", categoryId)
.filterTerm("status", status);
// ❌ 错误:使用 must 进行精确匹配会计算评分,浪费性能
// .term("categoryId", categoryId)
// .term("status", status)
return esService.search(queryBuilder, Program.class);
}
/**
* 2. 合理使用分页,避免深度分页
* from + size 超过 10000 会报错
*/
public PageResult<Program> safePageQuery(int pageNum, int pageSize) {
// 计算深度
int depth = (pageNum - 1) * pageSize + pageSize;
if (depth > 10000) {
// 深度分页使用 search_after
log.warn("分页深度超过10000,建议使用 searchAfter");
throw new IllegalArgumentException("分页深度超出限制,请使用 searchAfter");
}
EsQueryBuilder queryBuilder = EsQueryBuilder.builder(INDEX_NAME)
.page(pageNum, pageSize)
.sort("id", SortOrder.ASC);
return esService.searchPage(queryBuilder, pageNum, pageSize, Program.class);
}
/**
* 3. 只返回需要的字段,减少网络传输
*/
public List<Program> selectiveFieldsQuery(String keyword) {
EsQueryBuilder queryBuilder = EsQueryBuilder.builder(INDEX_NAME)
.match("title", keyword)
// 只返回需要的字段
.includes("id", "title", "minPrice", "showTime")
// 排除大字段
.excludes("description", "coverImage");
return esService.search(queryBuilder, Program.class);
}
/**
* 4. 使用 bool 查询组合多个条件
*/
public List<Program> complexBoolQuery(String keyword, List<Long> categoryIds,
Double minPrice, Double maxPrice,
List<Integer> excludeStatus) {
EsQueryBuilder queryBuilder = EsQueryBuilder.builder(INDEX_NAME)
// must: 必须匹配(影响评分)
.match("title", keyword)
// filter: 必须匹配(不影响评分,可缓存)
.filterTerms("categoryId", categoryIds)
.filterRange("minPrice", minPrice, maxPrice)
// must_not: 必须不匹配
.mustNotTerms("status", excludeStatus)
// should: 可选匹配(有则加分)
.should(QueryBuilders.matchQuery("actor", keyword))
.minimumShouldMatch(0);
return esService.search(queryBuilder, Program.class);
}
/**
* 5. 使用 constant_score 包装不需要评分的查询
*/
public List<Program> constantScoreQuery(Long categoryId) {
// 使用原生 QueryBuilder 实现 constant_score
QueryBuilder query = QueryBuilders.constantScoreQuery(
QueryBuilders.termQuery("categoryId", categoryId)
).boost(1.0f);
return esService.searchByQueryBuilder(INDEX_NAME, query, 100, Program.class);
}
/**
* 6. 批量查询优化:使用 multi_search
*/
public Map<String, List<Program>> multiSearch(List<Long> categoryIds) {
// 注意:需要扩展 ElasticsearchService 支持 multi_search
// 这里展示思路
Map<String, List<Program>> result = new HashMap<>();
for (Long categoryId : categoryIds) {
EsQueryBuilder queryBuilder = EsQueryBuilder.builder(INDEX_NAME)
.filterTerm("categoryId", categoryId)
.fromSize(0, 10);
List<Program> programs = esService.search(queryBuilder, Program.class);
result.put(String.valueOf(categoryId), programs);
}
return result;
}
/**
* 7. 高亮优化:限制高亮片段
*/
public List<SearchResult<Program>> optimizedHighlightQuery(String keyword) {
EsQueryBuilder queryBuilder = EsQueryBuilder.builder(INDEX_NAME)
.match("title", keyword)
.match("description", keyword);
// 自定义高亮配置
HighlightBuilder highlightBuilder = new HighlightBuilder()
.field(new HighlightBuilder.Field("title")
.fragmentSize(100) // 片段大小
.numOfFragments(1)) // 片段数量
.field(new HighlightBuilder.Field("description")
.fragmentSize(150)
.numOfFragments(2))
.preTags("<em>")
.postTags("</em>");
queryBuilder.setHighlightBuilder(highlightBuilder);
return esService.searchWithHighlight(queryBuilder, Program.class);
}
}
写入优化
/**
* 写入优化最佳实践
*/
@Slf4j
@Service
public class WriteOptimizationService {
@Autowired
private ElasticsearchDocumentService documentService;
@Autowired
private ElasticsearchIndexService indexService;
private static final String INDEX_NAME = "program";
/**
* 1. 批量写入优化
* - 合理设置批次大小(通常 1000-5000)
* - 控制并发写入线程数
*/
public void optimizedBatchWrite(List<Program> programs) {
int batchSize = 1000;
int totalSize = programs.size();
log.info("开始批量写入,总数: {}", totalSize);
long startTime = System.currentTimeMillis();
// 分批处理
for (int i = 0; i < totalSize; i += batchSize) {
int endIndex = Math.min(i + batchSize, totalSize);
List<Program> batch = programs.subList(i, endIndex);
BulkResult result = documentService.batchAdd(INDEX_NAME, batch);
if (result.hasFailures()) {
log.warn("批次 {}-{} 部分失败,失败数: {}",
i, endIndex, result.getFailureCount());
}
log.info("写入进度: {}/{}", endIndex, totalSize);
}
long elapsed = System.currentTimeMillis() - startTime;
log.info("批量写入完成,耗时: {}ms,速率: {} docs/s",
elapsed, totalSize * 1000L / elapsed);
}
/**
* 2. 大批量数据导入优化
* - 临时关闭副本
* - 增大 refresh_interval
* - 导入完成后恢复
*/
public void bulkImportWithOptimization(List<Program> programs) {
try {
// 优化索引设置
optimizeIndexForBulkImport();
// 执行批量导入
optimizedBatchWrite(programs);
} finally {
// 恢复索引设置
restoreIndexSettings();
}
}
private void optimizeIndexForBulkImport() {
Map<String, Object> settings = new HashMap<>();
settings.put("index.refresh_interval", "-1"); // 关闭自动刷新
settings.put("index.number_of_replicas", 0); // 临时关闭副本
indexService.updateSettings(INDEX_NAME, settings);
log.info("索引设置已优化");
}
private void restoreIndexSettings() {
Map<String, Object> settings = new HashMap<>();
settings.put("index.refresh_interval", "1s"); // 恢复刷新间隔
settings.put("index.number_of_replicas", 1); // 恢复副本
indexService.updateSettings(INDEX_NAME, settings);
// 强制刷新
indexService.refreshIndex(INDEX_NAME);
// 强制合并段
indexService.forceMerge(INDEX_NAME, 1);
log.info("索引设置已恢复");
}
/**
* 3. 使用异步写入
*/
@Async
public CompletableFuture<BulkResult> asyncBatchWrite(List<Program> programs) {
BulkResult result = documentService.batchAdd(INDEX_NAME, programs);
return CompletableFuture.completedFuture(result);
}
/**
* 4. 并行批量写入
*/
public void parallelBatchWrite(List<Program> programs, int threadCount) {
int batchSize = programs.size() / threadCount;
ExecutorService executor = Executors.newFixedThreadPool(threadCount);
List<Future<BulkResult>> futures = new ArrayList<>();
try {
for (int i = 0; i < threadCount; i++) {
int start = i * batchSize;
int end = (i == threadCount - 1) ? programs.size() : start + batchSize;
List<Program> batch = programs.subList(start, end);
Future<BulkResult> future = executor.submit(() ->
documentService.batchAdd(INDEX_NAME, batch));
futures.add(future);
}
// 等待所有任务完成
int totalSuccess = 0;
int totalFailure = 0;
for (Future<BulkResult> future : futures) {
BulkResult result = future.get();
totalSuccess += result.getSuccessCount();
totalFailure += result.getFailureCount();
}
log.info("并行写入完成,成功: {}, 失败: {}", totalSuccess, totalFailure);
} catch (Exception e) {
log.error("并行写入失败", e);
throw new RuntimeException(e);
} finally {
executor.shutdown();
}
}
}
缓存策略
/**
* ES 查询缓存策略
*/
@Slf4j
@Service
public class EsCacheService {
@Autowired
private ElasticsearchService esService;
@Autowired
private RedisTemplate<String, Object> redisTemplate;
private static final String CACHE_PREFIX = "es:program:";
private static final Duration CACHE_TTL = Duration.ofMinutes(5);
/**
* 1. 使用 Redis 缓存查询结果
*/
@SuppressWarnings("unchecked")
public List<Program> searchWithCache(String keyword) {
String cacheKey = CACHE_PREFIX + "search:" + DigestUtils.md5Hex(keyword);
// 先查缓存
Object cached = redisTemplate.opsForValue().get(cacheKey);
if (cached != null) {
log.debug("命中缓存: {}", cacheKey);
return (List<Program>) cached;
}
// 查询 ES
EsQueryBuilder queryBuilder = EsQueryBuilder.builder("program")
.match("title", keyword);
List<Program> result = esService.search(queryBuilder, Program.class);
// 写入缓存
if (!result.isEmpty()) {
redisTemplate.opsForValue().set(cacheKey, result, CACHE_TTL);
}
return result;
}
/**
* 2. 热门搜索词预热缓存
*/
@Scheduled(fixedRate = 300000) // 每5分钟
public void warmUpCache() {
List<String> hotKeywords = getHotKeywords();
for (String keyword : hotKeywords) {
try {
searchWithCache(keyword);
} catch (Exception e) {
log.warn("预热缓存失败: {}", keyword, e);
}
}
log.info("缓存预热完成,预热 {} 个关键词", hotKeywords.size());
}
/**
* 3. 数据更新时清除相关缓存
*/
public void invalidateCache(Long programId) {
// 清除该节目相关的所有缓存
Set<String> keys = redisTemplate.keys(CACHE_PREFIX + "*");
if (keys != null && !keys.isEmpty()) {
redisTemplate.delete(keys);
log.info("清除 {} 个缓存", keys.size());
}
}
/**
* 4. 分类统计结果缓存(变化不频繁的数据)
*/
@Cacheable(value = "es:stats", key = "'category'", unless = "#result == null")
public Map<String, Long> getCategoryStats() {
EsQueryBuilder queryBuilder = EsQueryBuilder.builder("program")
.termsAggregation("category_count", "categoryId", 100);
Map<String, AggregationResult> aggResult =
esService.searchAggregation(queryBuilder);
// 解析结果
Map<String, Long> stats = new HashMap<>();
AggregationResult categoryAgg = aggResult.get("category_count");
if (categoryAgg != null && categoryAgg.getBuckets() != null) {
for (AggregationResult.BucketData bucket : categoryAgg.getBuckets()) {
stats.put(bucket.getKey(), bucket.getDocCount());
}
}
return stats;
}
private List<String> getHotKeywords() {
// 从统计服务获取热门搜索词
return Arrays.asList("演唱会", "话剧", "音乐剧", "相声", "脱口秀");
}
}
监控与告警
/**
* ES 监控服务
*/
@Slf4j
@Service
public class EsMonitorService {
@Autowired
private RestHighLevelClient client;
@Autowired
private MeterRegistry meterRegistry;
/**
* 1. 记录查询耗时指标
*/
public <T> List<T> searchWithMetrics(EsQueryBuilder queryBuilder, Class<T> clazz) {
Timer.Sample sample = Timer.start(meterRegistry);
try {
// 执行查询
List<T> result = esService.search(queryBuilder, clazz);
// 记录成功指标
sample.stop(Timer.builder("es.query.duration")
.tag("index", queryBuilder.getIndexName())
.tag("status", "success")
.register(meterRegistry));
// 记录结果数量
meterRegistry.gauge("es.query.result.size",
Tags.of("index", queryBuilder.getIndexName()),
result.size());
return result;
} catch (Exception e) {
// 记录失败指标
sample.stop(Timer.builder("es.query.duration")
.tag("index", queryBuilder.getIndexName())
.tag("status", "error")
.register(meterRegistry));
meterRegistry.counter("es.query.error",
Tags.of("index", queryBuilder.getIndexName())).increment();
throw e;
}
}
/**
* 2. 集群健康检查
*/
@Scheduled(fixedRate = 60000)
public void checkClusterHealth() {
try {
ClusterHealthRequest request = new ClusterHealthRequest();
ClusterHealthResponse response = client.cluster()
.health(request, RequestOptions.DEFAULT);
String status = response.getStatus().name();
int numberOfNodes = response.getNumberOfNodes();
int activeShards = response.getActiveShards();
// 记录指标
meterRegistry.gauge("es.cluster.nodes", numberOfNodes);
meterRegistry.gauge("es.cluster.shards.active", activeShards);
// 状态告警
if ("RED".equals(status)) {
sendAlert("ES 集群状态异常: RED",
String.format("节点数: %d, 活跃分片: %d",
numberOfNodes, activeShards));
} else if ("YELLOW".equals(status)) {
log.warn("ES 集群状态: YELLOW,可能有副本未分配");
}
log.info("ES 集群健康检查: status={}, nodes={}, activeShards={}",
status, numberOfNodes, activeShards);
} catch (Exception e) {
log.error("ES 集群健康检查失败", e);
sendAlert("ES 集群健康检查失败", e.getMessage());
}
}
/**
* 3. 索引状态监控
*/
@Scheduled(fixedRate = 300000)
public void checkIndexStats() {
try {
IndicesStatsRequest request = new IndicesStatsRequest();
IndicesStatsResponse response = client.indices()
.stats(request, RequestOptions.DEFAULT);
response.getIndices().forEach((indexName, stats) -> {
long docCount = stats.getPrimaries().getDocs().getCount();
long storeSizeBytes = stats.getPrimaries().getStore().getSizeInBytes();
meterRegistry.gauge("es.index.docs.count",
Tags.of("index", indexName), docCount);
meterRegistry.gauge("es.index.store.size.bytes",
Tags.of("index", indexName), storeSizeBytes);
log.debug("索引 [{}] 统计: 文档数={}, 存储大小={}MB",
indexName, docCount, storeSizeBytes / 1024 / 1024);
});
} catch (Exception e) {
log.error("索引状态监控失败", e);
}
}
/**
* 4. 慢查询日志收集
*/
public void logSlowQuery(String indexName, String dsl, long elapsed) {
if (elapsed > 3000) { // 超过3秒
log.warn("[慢查询] 索引: {}, 耗时: {}ms, DSL: {}",
indexName, elapsed, dsl);
// 发送告警
if (elapsed > 10000) { // 超过10秒
sendAlert("ES 慢查询告警",
String.format("索引: %s, 耗时: %dms", indexName, elapsed));
}
}
}
private void sendAlert(String title, String content) {
// 发送告警通知(钉钉、企业微信、邮件等)
log.error("【告警】{}: {}", title, content);
}
}
十一、总结
💬 评论